diff --git a/apps/aevatar-console-web/src/shared/studio/document.test.ts b/apps/aevatar-console-web/src/shared/studio/document.test.ts index 61ed9b1a5f..e9c9214c53 100644 --- a/apps/aevatar-console-web/src/shared/studio/document.test.ts +++ b/apps/aevatar-console-web/src/shared/studio/document.test.ts @@ -45,6 +45,42 @@ describe('studio document helpers', () => { ); }); + it('inserts http_request with direct HTTP defaults and no role binding', () => { + const document: StudioWorkflowDocument = { + name: 'workspace-demo', + roles: [{ id: 'assistant' }], + steps: [], + }; + + const result = insertStepByType(document, 'http_request', { + targetRoleId: 'assistant', + }); + + expect(result.nodeId).toBe('step:http_request_step'); + expect(result.document.steps?.[0]).toEqual( + expect.objectContaining({ + id: 'http_request_step', + type: 'http_request', + originalType: 'http_request', + targetRole: null, + parameters: { + authentication: { scheme: 'bearer', secret_ref: '' }, + body: '', + body_mode: 'none', + headers: {}, + max_request_bytes: '65536', + max_response_bytes: '65536', + method: 'GET', + on_error: 'fail', + query: {}, + retry: '0', + timeout_ms: '30000', + url: '', + }, + }), + ); + }); + it('updates role fields and rewrites step role bindings', () => { const document: StudioWorkflowDocument = { name: 'workspace-demo', diff --git a/apps/aevatar-console-web/src/shared/studio/document.ts b/apps/aevatar-console-web/src/shared/studio/document.ts index 48fa78f538..e13fee79b8 100644 --- a/apps/aevatar-console-web/src/shared/studio/document.ts +++ b/apps/aevatar-console-web/src/shared/studio/document.ts @@ -80,6 +80,20 @@ const DEFAULT_PARAMETERS_BY_STEP_TYPE: Record> = retry: '0', on_error: 'fail', }, + http_request: { + method: 'GET', + url: '', + query: {}, + headers: {}, + body_mode: 'none', + body: '', + authentication: { scheme: 'bearer', secret_ref: '' }, + timeout_ms: '30000', + max_request_bytes: '65536', + max_response_bytes: '65536', + retry: '0', + on_error: 'fail', + }, emit: { event_type: 'workflow.completed', payload: '$input' }, human_input: { prompt: 'Please provide the missing input.', @@ -97,6 +111,7 @@ const DEFAULT_STEP_ID_BASE_BY_STEP_TYPE: Record = { dynamic_workflow: 'dynamic_workflow_step', human_approval: 'approval_step', human_input: 'human_input_step', + http_request: 'http_request_step', llm_call: 'llm_step', map_reduce: 'map_reduce_step', retrieve_facts: 'retrieve_facts_step', diff --git a/apps/aevatar-console-web/src/shared/studio/graph.ts b/apps/aevatar-console-web/src/shared/studio/graph.ts index fb7ef1ad4f..eeff5f53ff 100644 --- a/apps/aevatar-console-web/src/shared/studio/graph.ts +++ b/apps/aevatar-console-web/src/shared/studio/graph.ts @@ -126,7 +126,7 @@ export const STUDIO_GRAPH_CATEGORIES: readonly StudioGraphPrimitiveCategory[] = key: 'integration', label: 'Integration', color: '#10B981', - items: ['connector_call', 'emit'], + items: ['connector_call', 'http_request', 'emit'], }, { key: 'human', @@ -166,6 +166,7 @@ const STUDIO_STEP_TYPE_LABELS: Record = { guard: 'Guard', human_approval: 'Human approval', human_input: 'Human input', + http_request: 'HTTP Request', llm_call: 'LLM call', map_reduce: 'Map reduce', parallel: 'Parallel', diff --git a/apps/aevatar-console-web/src/shared/studio/nodeConfigFieldSchemas.ts b/apps/aevatar-console-web/src/shared/studio/nodeConfigFieldSchemas.ts index 12dc465c41..93d64fea4d 100644 --- a/apps/aevatar-console-web/src/shared/studio/nodeConfigFieldSchemas.ts +++ b/apps/aevatar-console-web/src/shared/studio/nodeConfigFieldSchemas.ts @@ -254,6 +254,13 @@ const STEP_TYPE_OPTIONS = [ ), value: 'connector_call', }, + { + label: message( + 'shared.studio.nodeConfiguration.stepType.option.httpRequest', + 'HTTP Request', + ), + value: 'http_request', + }, { label: message( 'shared.studio.nodeConfiguration.stepType.option.emit', @@ -523,6 +530,205 @@ const SCHEMAS_BY_STEP_TYPE: Record { }); }); + it('presents http_request as a first-class direct HTTP node with secret-backed auth fields', () => { + const schema = getStudioNodeConfigurationSchema('http_request'); + const values = readStudioNodeConfigurationValues('http_request', { + method: 'GET', + url: 'https://api.example.com/q1000', + authentication: { scheme: 'bearer', secret_ref: 'q1000-token' }, + timeout_ms: '20000', + max_request_bytes: '4096', + max_response_bytes: '65536', + }); + + expect(schema.fields.map((field) => field.parameterName)).toEqual( + expect.arrayContaining([ + 'method', + 'url', + 'query', + 'headers', + 'body_mode', + 'body', + 'authentication', + 'timeout_ms', + 'max_request_bytes', + 'max_response_bytes', + 'retry', + 'on_error', + ]), + ); + expect(schema.fields.find((field) => field.name === 'url')).toEqual( + expect.objectContaining({ + label: expect.objectContaining({ defaultMessage: 'URL' }), + required: true, + }), + ); + expect(schema.fields.find((field) => field.name === 'authentication')).toEqual( + expect.objectContaining({ + kind: 'object', + parameterName: 'authentication', + }), + ); + expect(values).toEqual( + expect.objectContaining({ + authentication: '{\n "scheme": "bearer",\n "secret_ref": "q1000-token"\n}', + maxResponseBytes: '65536', + method: 'GET', + timeoutMs: '20000', + url: 'https://api.example.com/q1000', + }), + ); + + expect( + applyStudioNodeConfigurationValues( + 'http_request', + { method: 'GET', url: 'https://api.example.com/q1000' }, + { + authentication: '{ "scheme": "bearer", "secret_ref": "q1000-token" }', + body: '', + bodyMode: 'none', + headers: '{}', + maxRequestBytes: '4096', + maxResponseBytes: '65536', + method: 'POST', + onError: 'fail', + query: '{}', + retry: '1', + timeoutMs: '20000', + url: 'https://api.example.com/q1000', + }, + ), + ).toEqual({ + authentication: { scheme: 'bearer', secret_ref: 'q1000-token' }, + body_mode: 'none', + headers: {}, + max_request_bytes: 4096, + max_response_bytes: 65536, + method: 'POST', + on_error: 'fail', + query: {}, + retry: 1, + timeout_ms: 20000, + url: 'https://api.example.com/q1000', + }); + }); + it('presents llm_call prompt_prefix as an Instruction field while preserving runtime parameters', () => { const schema = getStudioNodeConfigurationSchema('llm_call'); const values = readStudioNodeConfigurationValues('llm_call', { diff --git a/apps/aevatar-console-web/src/shared/studio/nodeConfigFields.ts b/apps/aevatar-console-web/src/shared/studio/nodeConfigFields.ts index 849a6020ce..5179523739 100644 --- a/apps/aevatar-console-web/src/shared/studio/nodeConfigFields.ts +++ b/apps/aevatar-console-web/src/shared/studio/nodeConfigFields.ts @@ -181,6 +181,114 @@ const CONNECTOR_CALL_FIELDS: readonly NodeConfigFieldSource[] = [ }, ]; +const HTTP_REQUEST_FIELDS: readonly NodeConfigFieldSource[] = [ + { + name: 'method', + label: message('shared.studio.nodeConfigFields.httpRequest.method.label', 'Method'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.method.description', + 'HTTP method for this direct request.', + ), + enumValues: ['GET', 'POST', 'PUT', 'PATCH', 'DELETE'], + type: 'string', + default: 'GET', + }, + { + name: 'url', + label: message('shared.studio.nodeConfigFields.httpRequest.url.label', 'URL'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.url.description', + 'Absolute HTTPS URL called by this workflow step.', + ), + required: true, + type: 'string', + }, + { + name: 'query', + label: message('shared.studio.nodeConfigFields.httpRequest.query.label', 'Query parameters'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.query.description', + 'Query parameter object appended to the URL.', + ), + type: 'object', + }, + { + name: 'headers', + label: message('shared.studio.nodeConfigFields.httpRequest.headers.label', 'Headers'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.headers.description', + 'Request headers. Use authentication.secret_ref instead of raw Authorization values.', + ), + type: 'object', + }, + { + name: 'authentication', + label: message('shared.studio.nodeConfigFields.httpRequest.authentication.label', 'Authentication'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.authentication.description', + 'Object such as {"scheme":"bearer","secret_ref":"saved-secret"}.', + ), + type: 'object', + }, + { + name: 'timeout_ms', + label: message('shared.studio.nodeConfigFields.httpRequest.timeoutMs.label', 'Timeout ms'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.timeoutMs.description', + 'Direct HTTP request timeout in milliseconds.', + ), + type: 'number', + default: '30000', + }, + { + name: 'max_response_bytes', + label: message( + 'shared.studio.nodeConfigFields.httpRequest.maxResponseBytes.label', + 'Max response bytes', + ), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.maxResponseBytes.description', + 'Maximum response body size accepted by the workflow.', + ), + type: 'number', + default: '65536', + }, + { + name: 'max_request_bytes', + label: message( + 'shared.studio.nodeConfigFields.httpRequest.maxRequestBytes.label', + 'Max request bytes', + ), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.maxRequestBytes.description', + 'Maximum request body size sent by the workflow.', + ), + type: 'number', + default: '65536', + }, + { + name: 'retry', + label: message('shared.studio.nodeConfigFields.httpRequest.retry.label', 'Retry'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.retry.description', + 'Retry count for failed HTTP attempts.', + ), + type: 'number', + default: '0', + }, + { + name: 'on_error', + label: message('shared.studio.nodeConfigFields.httpRequest.onError.label', 'On error'), + description: message( + 'shared.studio.nodeConfigFields.httpRequest.onError.description', + 'Failure behavior when the HTTP request cannot complete.', + ), + enumValues: ['fail', 'continue'], + type: 'string', + default: 'fail', + }, +]; + const LLM_CALL_FIELDS: readonly NodeConfigFieldSource[] = [ { name: PROMPT_PREFIX_PARAMETER, @@ -479,6 +587,8 @@ function createKnownFieldSources( switch (normalizeStepType(stepType)) { case 'connector_call': return [...CONNECTOR_CALL_FIELDS]; + case 'http_request': + return [...HTTP_REQUEST_FIELDS]; case LLM_CALL_STEP_TYPE: return [...LLM_CALL_FIELDS]; default: diff --git a/docs/canon/connector.md b/docs/canon/connector.md index 4c6ffb826a..935d1ddaa9 100644 --- a/docs/canon/connector.md +++ b/docs/canon/connector.md @@ -18,6 +18,7 @@ owner: eanzhao - `IConnector` / `IConnectorRegistry` 是统一外部调用契约,定义在 `Aevatar.Foundation.Abstractions`。 - Connector 的定义是中心化配置(`~/.aevatar/connectors.json`)。 - Workflow 里通过 `type: connector_call` + `parameters.connector` 使用命名 Connector。 +- Direct workflow HTTP uses `type: http_request`; it does not read `parameters.connector`, does not consume `connectors.json`, and does not require `IConnectorRegistry`. - 角色(role)里的 `connectors` 是授权白名单,不是连接定义本身。 简化链路: @@ -116,6 +117,8 @@ Builder 当前行为: ## 3. Workflow/Agent 如何使用 Connector +This section describes named connectors. Use `http_request` when one workflow owns the exact URL, method, headers, body, limits, and `authentication.secret_ref`. Use a named HTTP connector when the integration should be reusable, centrally governed, role-allowlisted, or constrained by connector-level base URL and path policy. Both paths delegate HTTP transport to `IOutboundHttpRequestExecutor`. + ## 3.1 Workflow YAML 角色授权 在 workflow `roles` 中: @@ -140,7 +143,7 @@ roles: ## 3.2 connector_call 执行主链路 -`ConnectorCallModule` 处理 `StepRequestEvent`(`step_type == connector_call`): +`ConnectorCallModule` 处理 named connector `StepRequestEvent`(`step_type == connector_call` or `secure_connector_call`): 1. 读取参数: - `connector`(或 `connector_name`)必填 @@ -151,17 +154,23 @@ roles: 4. 构造 `ConnectorRequest` 并调用 `IConnector.ExecuteAsync()`; 5. 根据结果发布 `StepCompletedEvent`。 -Ergonomic 别名(解析期归一化到 `connector_call`): +Ergonomic aliases: -- `http_get` / `http_post` / `http_put` / `http_delete` - - 归一化后仍是 `connector_call` - - 若未显式指定 `method`,会自动补 `GET/POST/PUT/DELETE` - `mcp_call` - - 归一化后仍是 `connector_call` + - Normalizes to `connector_call`. - 仅写 `tool` 且未写 `operation/action` 时,会自动补 `operation=` - `cli_call` - - 归一化后仍是 `connector_call` + - Normalizes to `connector_call`. - 不改变执行语义(仍通过命名 connector 调用) +- `bridge_call` + - Normalizes to `connector_call`. + +HTTP convenience aliases are not named connector aliases anymore: + +- `http_get` / `http_post` / `http_put` / `http_delete` + - Normalize to `http_request`. + - Add `method=GET/POST/PUT/DELETE` if no method is explicitly configured. + - Require direct HTTP fields such as `url`; they do not require `parameters.connector`. 容错语义: @@ -233,13 +242,14 @@ Actor audit facts. Recovery then revokes both protected request and completion m ### HTTP Connector -- 方法默认 `POST`,可由参数 `method` 覆盖; +- Named HTTP Connector 方法默认 `POST`,可由参数 `method` 覆盖; - 路径优先 `operation`,其次参数 `path`; - 必须通过 `allowedMethods/allowedPaths` 白名单; - 强制校验目标 URL 不能逃逸 `baseUrl` 的 scheme/host/port; - 可用 `allowedInputKeys` 校验 payload JSON key; - 返回 HTTP 状态和耗时元数据。 - `defaultHeaders` 只用于非 secret 静态 header。secret-bearing header 必须使用 `auth.type=secret_ref_header`,避免 raw secret 被复制进 connector config、workflow 参数、annotations、read model 或通用 bag。 +- Named HTTP Connector and direct `http_request` share the same hardened outbound transport executor. The connector adds reusable base URL, path, method, input-key, and connector credential policy before delegating transport, where request and response byte limits are enforced. `secret_ref_header` 配置形状: diff --git a/docs/canon/workflow-primitives.md b/docs/canon/workflow-primitives.md index 514cfeecec..c96859ec39 100644 --- a/docs/canon/workflow-primitives.md +++ b/docs/canon/workflow-primitives.md @@ -95,7 +95,7 @@ steps: Runtime semantics: -- `tool_call`, `connector_call`, and `secure_connector_call` are the v1.1 side-effecting primitive set. When one declares `compensation`, dispatch first records a `PROVISIONAL` ledger entry before the external side-effect boundary. +- `tool_call`, `http_request`, `connector_call`, and `secure_connector_call` are the v1.1 side-effecting primitive set. When one declares `compensation`, dispatch first records a `PROVISIONAL` ledger entry before the external side-effect boundary. - A successful completion confirms a matching provisional entry as `CONFIRMED` and fills captured output. If no dispatch event exists, legacy success still appends one `CONFIRMED` entry. - A callee-confirmed failure removes the matching provisional entry. Timeout, force-fail, or stop-to-failure paths set `failure_outcome = OUTCOME_UNCERTAIN`, keep the provisional entry, and let compensation treat undoing a not-applied side effect as a safe no-op. - If a later terminal failure occurs while the ledger is non-empty, compensation runs in reverse ledger order over both `PROVISIONAL` and `CONFIRMED` entries. @@ -674,9 +674,60 @@ steps: ## 6. Integration 原语 -### `connector_call`(别名:`bridge_call`、`cli_call`、`mcp_call`、`http_get`、`http_post`、`http_put`、`http_delete`) +### `http_request` (aliases: `http_get`, `http_post`, `http_put`, `http_delete`) + +- Purpose: send one direct outbound HTTP request from the workflow without looking up a named connector in `IConnectorRegistry`. +- Common parameters: `method`, `url`, `query`, `headers`, `body_mode`, `body`, `authentication`, `timeout_ms`, `max_request_bytes`, `max_response_bytes`, `max_redirects`, `retry`, `on_error`. +- Stable control semantics are typed as `WorkflowHttpRequestOptions` and `WorkflowHttpRequestAuthentication`, then carried through `WorkflowStepParameters.http_request`. They are not inferred from a generic connector name or arbitrary response JSON. +- `http_request` side effects are at-least-once. The workflow actor uses the logical run id + step id + logical attempt to resolve and persist the typed `idempotency_key`; pending replay and physical retry reuse that key and the executor sends it as `Idempotency-Key` when non-empty. This key is for callee-side dedup only; the workflow engine does not promise exactly-once. +- Authentication must use `authentication.secret_ref` or compatible `credential_ref` input names. Raw secret values and raw `Authorization` headers are rejected before dispatch. The raw secret is resolved through `ICredentialProvider` only at execution time, and secret-derived values are redacted from output, error text, annotations, read models, traces, and UI-visible metadata. +- Supported authentication schemes are `bearer`, `header`, and `secret_ref_header`. `bearer` produces an HTTP `Authorization: Bearer ` header inside the executor boundary; `header` and `secret_ref_header` require `authentication.header_name` and may use `authentication.header_value_prefix`. +- HTTPS is required by default. `allow_insecure_http` is a development-only override and does not allow private network destinations. +- The shared outbound HTTP executor validates every target and redirect before dispatch, disables automatic redirects, enforces a bounded redirect count, and validates DNS again at the socket connection boundary to defend against rebinding. It rejects loopback, link-local, multicast, cloud metadata style, private, and carrier-grade NAT destinations unless explicitly allowed by a named connector boundary, and applies timeout plus maximum request and response byte limits. +- Failure semantics are explicit: DNS resolution failure, blocked egress target, redirect policy denial, authentication failure, timeout, non-success HTTP status, invalid URL, and response-size overflow all fail the step unless `on_error: continue` is set. +- Direct `http_request` and named HTTP connectors share the same hardened `IOutboundHttpRequestExecutor`. The authoring models differ: `http_request` owns the concrete URL/method/options in workflow YAML, while a named HTTP connector owns reusable base URL, path allowlists, and centralized governance in `connectors.json` or a future Team Connector catalog. +- Ergonomic aliases normalize to `http_request`: + - `http_get`: adds `method=GET` if the method is absent. + - `http_post`: adds `method=POST` if the method is absent. + - `http_put`: adds `method=PUT` if the method is absent. + - `http_delete`: adds `method=DELETE` if the method is absent. + +```yaml +steps: + - id: fetch_snapshot + type: http_request + parameters: + method: GET + url: "https://api.example.com/q1000" + query: + source: "dashboard" + headers: + X-Api-Version: "1" + authentication: + scheme: bearer + secret_ref: q1000_bridge_token + timeout_ms: "20000" + max_request_bytes: "65536" + max_response_bytes: "65536" + max_redirects: "2" + retry: + max_attempts: 3 + backoff: exponential + delay_ms: 1000 +``` + +```yaml +steps: + - id: get_health + type: http_get + parameters: + url: "https://api.example.com/healthz" + timeout_ms: "5000" +``` + +### `connector_call`(别名:`bridge_call`、`cli_call`、`mcp_call`) -- 作用:调用外部 connector(HTTP/CLI/MCP 等),支持重试和降级策略。 +- 作用:调用外部 named connector(HTTP/CLI/MCP/host_callback 等),支持重试、命名治理、角色 allowlist 和降级策略。 - 常用参数:`connector`、`operation`、`retry`、`timeout_ms`、`optional`、`on_missing`、`on_error`。 - `connector_call` / `secure_connector_call` side effect 是 at-least-once。workflow actor 按 logical run id + step id + logical attempt 解析并持久化 typed `idempotency_key`;若 step 声明 `compensation`,同一 seam 先写入 `PROVISIONAL` compensation ledger,再发布 connector request。connector physical retry / pending replay 复用同一个 key;HTTP connector 会在 key 非空时发送 `Idempotency-Key` header,其他 connector 可按自身边界使用或忽略。该 key 不提供 engine-side dedup 或 exactly-once。 - `approval.policy: required` enables actor-owned durable approval coordination before connector dispatch. The step must provide `approval.service_ref`, `approval.node_id`, `approval.http_verb`, `approval.resource`, `approval.permission_scope`, `approval.expiration_seconds`, and a stable `idempotency_key`. `approval.status_check_interval_seconds` defaults to 2. @@ -684,10 +735,10 @@ steps: - Approval state survives restart through actor state plus durable self callbacks. NyxID submission or status uncertainty fails closed; an indeterminate submission is not retried because NyxID creates a unique request for each submission. - Approved execution revalidates the remote binding, action, digest, caller authority, scope, node, service, permission scope, and effective expiry immediately before dispatch. HTTP approvals also require the approved verb/resource to match the concrete connector method/path. Dispatch replay and connector retries reuse the same physical `idempotency_key`; approval success and connector success remain separate persisted facts. - The Actor persists a dispatch acknowledgement and keeps the exact pending `StepCompletedEvent` as protected Protobuf until publication succeeds. Restart recovery can therefore redispatch an unacknowledged invocation or republish an acknowledged external result without copying response content into audit facts or public read models. -- Ergonomic 说明(统一归一化到 `connector_call`): - - `http_get`/`http_post`/`http_put`/`http_delete`:自动补 `method=GET/POST/PUT/DELETE`(若未显式提供)。 - - `mcp_call`:若只写 `tool` 且未写 `operation/action`,会自动补 `operation=`。 - - `cli_call`:仅语义别名,不改变执行语义。 +- Ergonomic aliases: + - `mcp_call`: normalizes to `connector_call`; when only `tool` is present and `operation/action` is absent, parsing adds `operation=`. + - `cli_call`: normalizes to `connector_call`; execution still uses a named connector. + - `bridge_call`: normalizes to `connector_call`. ```yaml steps: @@ -702,16 +753,6 @@ steps: on_error: "continue" ``` -```yaml -steps: - - id: get_health - type: http_get - target_role: coordinator - parameters: - connector: "internal_http" - path: "/healthz" -``` - ```yaml steps: - id: create_resource diff --git a/docs/canon/workflow-runtime.md b/docs/canon/workflow-runtime.md index 32bab980d4..7b37417d7c 100644 --- a/docs/canon/workflow-runtime.md +++ b/docs/canon/workflow-runtime.md @@ -268,6 +268,7 @@ roles: | **引擎** | N/A | `WorkflowExecutionKernel` | 按步骤顺序派发,收到完成事件后推进下一步或结束 | | **执行** | `llm_call` | `LLMCallModule` | 向目标 RoleGAgent 发 `ChatRequestEvent`,等回复转 `StepCompletedEvent` | | | `tool_call` | `ToolCallModule` | 调用已注册的 Agent 工具(MCP/Skills) | +| | `http_request` | `ConnectorCallModule` | 执行 workflow-owned direct HTTP request,不要求 named connector lookup | | | `connector_call` | `ConnectorCallModule` | 按名称调用配置好的 HTTP/CLI/MCP/host_callback connector | | **并行** | `parallel` | `ParallelFanOutModule` | 拆 N 个子步骤并行发给不同 role,收齐后合并,可选触发 typed vote agreement | | **共识** | `vote` | `VoteAgreementModule` | 基于 typed candidate/rule/decision 做结构化 agreement 判定(`vote_consensus` 为别名) | @@ -339,7 +340,7 @@ POST /api/chat { prompt, workflow?, workflowYaml?, source? } │ ├── 对应模块处理 StepRequestEvent │ ├── LLMCallModule: 转 ChatRequestEvent → SendTo RoleGAgent → 等 TextMessageEndEvent → StepCompletedEvent - │ ├── ConnectorCallModule: 查 registry → 执行 connector → StepCompletedEvent + │ ├── ConnectorCallModule: `http_request` 走 typed direct HTTP;`connector_call` 查 registry → 执行 connector → StepCompletedEvent │ ├── ParallelFanOutModule: 拆子步骤 → 收齐合并 → 可选投票 → StepCompletedEvent │ └── ...其他模块同理 │ @@ -360,7 +361,7 @@ POST /api/chat { prompt, workflow?, workflowYaml?, source? } ### Saga 补偿生命周期 -Workflow step 可以通过 `compensation` 声明一个已存在的 step id。静态校验阶段会解析该目标,引用不存在的补偿步骤会被拒绝。运行时中,`tool_call`、`connector_call`、`secure_connector_call` 这三个 side-effecting primitive 在 dispatch 前会由 `WorkflowRunGAgent` 持久化 `CompensableStepDispatchedEvent`,先写入 `PROVISIONAL` ledger 项;其他 primitive 即使声明 compensation,也只在成功完成后按 legacy 路径写入 `CONFIRMED` ledger 项。`compensable_ledger` 归 `WorkflowRunGAgent` 持有,是 run actor 的权威状态。 +Workflow step 可以通过 `compensation` 声明一个已存在的 step id。静态校验阶段会解析该目标,引用不存在的补偿步骤会被拒绝。运行时中,`tool_call`、`http_request`、`connector_call`、`secure_connector_call` 这些 side-effecting primitive 在 dispatch 前会由 `WorkflowRunGAgent` 持久化 `CompensableStepDispatchedEvent`,先写入 `PROVISIONAL` ledger 项;其他 primitive 即使声明 compensation,也只在成功完成后按 legacy 路径写入 `CONFIRMED` ledger 项。`compensable_ledger` 归 `WorkflowRunGAgent` 持有,是 run actor 的权威状态。 成功完成会把匹配的 `PROVISIONAL` ledger 项确认成 `CONFIRMED` 并补齐 captured output;没有 provisional 项的 legacy success 仍追加一条 `CONFIRMED` ledger。失败完成通过 typed `WorkflowStepFailureOutcome` 对账:`CALLEE_CONFIRMED`(含默认 `UNSPECIFIED`)删除匹配 provisional,表示 callee 已确认没有可补偿副作用;`OUTCOME_UNCERTAIN` 保留 provisional,timeout / force-fail / stop-to-failure 这类中断按“副作用可能已发生”处理。 @@ -497,9 +498,37 @@ services.TryAddEnumerable(ServiceDescriptor.Singleton public sealed class HttpConnector : IConnector { - private static readonly HttpClient SharedHttpClient = new(); + private static readonly HttpClient SharedHttpClient = new( + new SocketsHttpHandler + { + AllowAutoRedirect = false, + }); private readonly HttpClient? _client; private readonly IHttpClientFactory? _httpClientFactory; private readonly string _httpClientName; private readonly IConnectorRequestAuthorizationProvider? _authorizationProvider; + private readonly IOutboundHttpRequestExecutor? _outboundHttpRequestExecutor; private readonly Uri _baseUri; private readonly HashSet _allowedMethods; private readonly string[] _allowedPathPatterns; @@ -35,7 +39,8 @@ public HttpConnector( IHttpClientFactory? httpClientFactory = null, string? httpClientName = null, IConnectorRequestAuthorizationProvider? authorizationProvider = null, - HttpClient? client = null) + HttpClient? client = null, + IOutboundHttpRequestExecutor? outboundHttpRequestExecutor = null) { if (string.IsNullOrWhiteSpace(name)) throw new ArgumentException("name is required", nameof(name)); if (string.IsNullOrWhiteSpace(baseUrl)) throw new ArgumentException("baseUrl is required", nameof(baseUrl)); @@ -57,6 +62,7 @@ public HttpConnector( _httpClientFactory = httpClientFactory; _httpClientName = string.IsNullOrWhiteSpace(httpClientName) ? Name : httpClientName.Trim(); _authorizationProvider = authorizationProvider; + _outboundHttpRequestExecutor = outboundHttpRequestExecutor; } /// @@ -132,72 +138,55 @@ public async Task ExecuteAsync(ConnectorRequest request, Canc }; } - var sw = Stopwatch.StartNew(); - using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(ct); - timeoutCts.CancelAfter(TimeSpan.FromMilliseconds(timeoutMs)); - try { - using var msg = new HttpRequestMessage(new HttpMethod(method), targetUri); - foreach (var (key, value) in _defaultHeaders) - msg.Headers.TryAddWithoutValidation(key, value); - - if (_authorizationProvider != null) - await _authorizationProvider.ApplyAsync(msg, timeoutCts.Token); - - ApplyRequestAuthorization(msg, request.HttpAuthorization); - ApplyIdempotencyKey(msg, request.IdempotencyKey); - - if (request.Parameters.TryGetValue("content_type", out var contentType) && - !string.IsNullOrWhiteSpace(contentType)) - { - msg.Content = new StringContent(request.Payload ?? "", Encoding.UTF8, contentType); - } - else if (method is not "GET" and not "HEAD") - { - msg.Content = new StringContent(request.Payload ?? "", Encoding.UTF8, "application/json"); - } - - if (!msg.Headers.Accept.Any()) - msg.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("application/json")); - - using var response = await ResolveClient().SendAsync(msg, timeoutCts.Token); - var body = await response.Content.ReadAsStringAsync(timeoutCts.Token); - sw.Stop(); - - return new ConnectorResponse - { - Success = response.IsSuccessStatusCode, - Output = body, - Error = response.IsSuccessStatusCode ? "" : BuildHttpErrorMessage(response, body), - Metadata = new Dictionary + var outboundHeaders = await BuildOutboundHeadersAsync(method, targetUri, request, ct); + var contentType = request.Parameters.TryGetValue("content_type", out var configuredContentType) && + !string.IsNullOrWhiteSpace(configuredContentType) + ? configuredContentType.Trim() + : "application/json"; + var maxResponseBytes = request.Parameters.TryGetValue("max_response_bytes", out var rawMaxResponseBytes) && + int.TryParse(rawMaxResponseBytes, out var parsedMaxResponseBytes) + ? parsedMaxResponseBytes + : DefaultOutboundHttpRequestExecutor.DefaultMaxResponseBytes; + var maxRequestBytes = request.Parameters.TryGetValue("max_request_bytes", out var rawMaxRequestBytes) && + int.TryParse(rawMaxRequestBytes, out var parsedMaxRequestBytes) + ? parsedMaxRequestBytes + : DefaultOutboundHttpRequestExecutor.DefaultMaxRequestBytes; + var maxRedirects = request.Parameters.TryGetValue("max_redirects", out var rawMaxRedirects) && + int.TryParse(rawMaxRedirects, out var parsedMaxRedirects) + ? parsedMaxRedirects + : DefaultOutboundHttpRequestExecutor.DefaultMaxRedirects; + + var response = await ResolveExecutor().ExecuteAsync( + new OutboundHttpRequest { - ["connector.http.status_code"] = ((int)response.StatusCode).ToString(), - ["connector.http.reason"] = response.ReasonPhrase ?? "", - ["connector.http.method"] = method, - ["connector.http.url"] = targetUri.ToString(), - ["connector.http.duration_ms"] = sw.Elapsed.TotalMilliseconds.ToString("F2"), + Method = method, + Url = targetUri.ToString(), + Headers = outboundHeaders.Headers, + Authorization = outboundHeaders.Authorization, + IdempotencyKey = request.IdempotencyKey, + Body = request.Payload ?? string.Empty, + ContentType = contentType, + TimeoutMs = timeoutMs, + MaxRequestBytes = maxRequestBytes, + MaxResponseBytes = maxResponseBytes, + MaxRedirects = maxRedirects, + AllowInsecureHttp = string.Equals(targetUri.Scheme, Uri.UriSchemeHttp, StringComparison.OrdinalIgnoreCase), + AllowPrivateNetwork = true, }, - }; - } - catch (OperationCanceledException) when (!ct.IsCancellationRequested) - { - sw.Stop(); + ct); + return new ConnectorResponse { - Success = false, - Error = $"http timeout after {timeoutMs}ms", - Metadata = new Dictionary - { - ["connector.http.method"] = method, - ["connector.http.url"] = targetUri.ToString(), - ["connector.http.duration_ms"] = sw.Elapsed.TotalMilliseconds.ToString("F2"), - }, + Success = response.Success, + Output = response.Output, + Error = response.Error, + Metadata = response.Metadata, }; } - catch (Exception ex) + catch (Exception ex) when (!ct.IsCancellationRequested) { - sw.Stop(); return new ConnectorResponse { Success = false, @@ -206,7 +195,6 @@ public async Task ExecuteAsync(ConnectorRequest request, Canc { ["connector.http.method"] = method, ["connector.http.url"] = targetUri.ToString(), - ["connector.http.duration_ms"] = sw.Elapsed.TotalMilliseconds.ToString("F2"), }, }; } @@ -294,6 +282,44 @@ private static string RawBodyPreview(string body) return trimmed.Length <= 200 ? trimmed : $"{trimmed[..200]}..."; } + private async Task<(Dictionary Headers, string Authorization)> BuildOutboundHeadersAsync( + string method, + Uri targetUri, + ConnectorRequest request, + CancellationToken ct) + { + using var message = new HttpRequestMessage(new HttpMethod(method), targetUri); + foreach (var (key, value) in _defaultHeaders) + message.Headers.TryAddWithoutValidation(key, value); + + if (_authorizationProvider != null) + await _authorizationProvider.ApplyAsync(message, ct); + + ApplyRequestAuthorization(message, request.HttpAuthorization); + + var headers = new Dictionary(StringComparer.OrdinalIgnoreCase); + foreach (var header in message.Headers) + { + if (string.Equals(header.Key, "Authorization", StringComparison.OrdinalIgnoreCase)) + continue; + + headers[header.Key] = string.Join(",", header.Value); + } + + return (headers, message.Headers.Authorization?.ToString() ?? string.Empty); + } + + private IOutboundHttpRequestExecutor ResolveExecutor() + { + if (_outboundHttpRequestExecutor != null) + return _outboundHttpRequestExecutor; + + if (_client == null && _httpClientFactory == null) + return new DefaultOutboundHttpRequestExecutor(); + + return new DefaultOutboundHttpRequestExecutor(ResolveClient()); + } + private HttpClient ResolveClient() { if (_httpClientFactory != null) diff --git a/src/Aevatar.Configuration/README.md b/src/Aevatar.Configuration/README.md index 27537cc446..692011d5b5 100644 --- a/src/Aevatar.Configuration/README.md +++ b/src/Aevatar.Configuration/README.md @@ -35,6 +35,8 @@ Connector 是框架提供的**命名外部调用抽象**:在认知工作流(Cognitive Workflow)中,用统一契约调用外部能力,而不必在 YAML 里写死 URL、命令或 MCP 细节。 +This README documents named connectors loaded from `connectors.json`. Direct one-off workflow HTTP is modeled separately as `type: http_request`; it stores typed request options in workflow YAML, resolves `authentication.secret_ref` only at execution time, and does not consume `connectors.json` or `IConnectorRegistry`. + - **使用场景**:由带 role 的 workflow(如 MAKER 分析)在某个步骤里按「名称」调用外部服务或本地命令。 - **谁消费**:工作流步骤类型 `connector_call`。步骤里通过 `parameters.connector` 指定已配置的 connector 名称,运行时从 `IConnectorRegistry` 解析并执行。 - **支持类型**:`http`(HTTP 接口)、`cli`(本地可执行命令)、`mcp`(MCP 服务器工具调用)。每种类型有独立的策略字段(如 baseUrl、command、allowedTools 等),用于安全与行为控制。 diff --git a/src/Aevatar.Foundation.Abstractions/Connectors/OutboundHttpContracts.cs b/src/Aevatar.Foundation.Abstractions/Connectors/OutboundHttpContracts.cs new file mode 100644 index 0000000000..8c86865b6d --- /dev/null +++ b/src/Aevatar.Foundation.Abstractions/Connectors/OutboundHttpContracts.cs @@ -0,0 +1,59 @@ +using System.Net; + +namespace Aevatar.Foundation.Abstractions.Connectors; + +public interface IOutboundHttpRequestExecutor +{ + Task ExecuteAsync( + OutboundHttpRequest request, + CancellationToken ct = default); +} + +public interface IOutboundHttpDnsResolver +{ + ValueTask> GetHostAddressesAsync( + string host, + CancellationToken ct = default); +} + +public sealed class OutboundHttpRequest +{ + public string Method { get; init; } = "GET"; + + public string Url { get; init; } = ""; + + public IReadOnlyDictionary Query { get; init; } = new Dictionary(); + + public IReadOnlyDictionary Headers { get; init; } = new Dictionary(); + + public string Authorization { get; init; } = ""; + + public string IdempotencyKey { get; init; } = ""; + + public string Body { get; init; } = ""; + + public string ContentType { get; init; } = ""; + + public int TimeoutMs { get; init; } + + public int MaxRequestBytes { get; init; } + + public int MaxResponseBytes { get; init; } + + public int MaxRedirects { get; init; } + + public bool AllowInsecureHttp { get; init; } + + public bool AllowPrivateNetwork { get; init; } +} + +public sealed class OutboundHttpResponse +{ + public bool Success { get; init; } + + public string Output { get; init; } = ""; + + public string Error { get; init; } = ""; + + public Dictionary Metadata { get; init; } = []; +} diff --git a/src/Aevatar.Foundation.Core/Connectors/DefaultOutboundHttpRequestExecutor.cs b/src/Aevatar.Foundation.Core/Connectors/DefaultOutboundHttpRequestExecutor.cs new file mode 100644 index 0000000000..965ad293ea --- /dev/null +++ b/src/Aevatar.Foundation.Core/Connectors/DefaultOutboundHttpRequestExecutor.cs @@ -0,0 +1,467 @@ +using System.Diagnostics; +using System.Net; +using System.Net.Http.Headers; +using System.Net.Sockets; +using System.Text; +using System.Text.Json; +using Aevatar.Foundation.Abstractions.Connectors; + +namespace Aevatar.Foundation.Core.Connectors; + +public sealed class DefaultOutboundHttpRequestExecutor : IOutboundHttpRequestExecutor +{ + public const int DefaultTimeoutMs = 30_000; + public const int DefaultMaxRequestBytes = 65_536; + public const int DefaultMaxResponseBytes = 65_536; + public const int DefaultMaxRedirects = 3; + + private static readonly HttpRequestOptionsKey AllowPrivateNetworkOption = + new("Aevatar.Foundation.Connectors.AllowPrivateNetwork"); + + private static readonly HttpClient SharedHttpClient = + CreateHardenedHttpClient(new DefaultOutboundHttpDnsResolver()); + + private readonly HttpClient _client; + private readonly IOutboundHttpDnsResolver _dnsResolver; + + public DefaultOutboundHttpRequestExecutor() + : this(SharedHttpClient) + { + } + + public DefaultOutboundHttpRequestExecutor(IOutboundHttpDnsResolver dnsResolver) + : this(CreateHardenedHttpClient(dnsResolver ?? throw new ArgumentNullException(nameof(dnsResolver))), dnsResolver) + { + } + + public DefaultOutboundHttpRequestExecutor( + HttpClient client, + IOutboundHttpDnsResolver? dnsResolver = null) + { + _client = client ?? throw new ArgumentNullException(nameof(client)); + _dnsResolver = dnsResolver ?? new DefaultOutboundHttpDnsResolver(); + } + + public async Task ExecuteAsync( + OutboundHttpRequest request, + CancellationToken ct = default) + { + ArgumentNullException.ThrowIfNull(request); + + if (!Uri.TryCreate(request.Url?.Trim(), UriKind.Absolute, out var targetUri)) + return Failure("http_request url must be an absolute URL"); + + var method = NormalizeMethod(request.Method); + var timeoutMs = ClampOrDefault(request.TimeoutMs, 100, 300_000, DefaultTimeoutMs); + var maxRequestBytes = ClampOrDefault(request.MaxRequestBytes, 1, 10 * 1024 * 1024, DefaultMaxRequestBytes); + var maxResponseBytes = ClampOrDefault(request.MaxResponseBytes, 1, 10 * 1024 * 1024, DefaultMaxResponseBytes); + var maxRedirects = ClampOrDefault(request.MaxRedirects, 0, 10, DefaultMaxRedirects); + var currentUri = ApplyQuery(targetUri, request.Query); + var currentMethod = method; + var body = request.Body ?? string.Empty; + var contentType = string.IsNullOrWhiteSpace(request.ContentType) + ? "application/json" + : request.ContentType.Trim(); + var sw = Stopwatch.StartNew(); + + if (Encoding.UTF8.GetByteCount(body) > maxRequestBytes) + return Failure($"http request exceeded {maxRequestBytes} bytes", currentMethod, currentUri, sw); + + using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(ct); + timeoutCts.CancelAfter(TimeSpan.FromMilliseconds(timeoutMs)); + + for (var redirectCount = 0; ; redirectCount++) + { + var policyError = await ValidateTargetAsync(currentUri, request, timeoutCts.Token); + if (!string.IsNullOrWhiteSpace(policyError)) + return Failure(policyError, currentMethod, currentUri, sw); + + using var message = BuildRequestMessage( + currentMethod, + currentUri, + request, + body, + contentType); + + try + { + using var response = await _client.SendAsync( + message, + HttpCompletionOption.ResponseHeadersRead, + timeoutCts.Token); + + if (IsRedirect(response.StatusCode)) + { + if (redirectCount >= maxRedirects) + return Failure("http redirect limit exceeded", currentMethod, currentUri, sw); + + if (response.Headers.Location == null) + return Failure("http redirect response missing Location header", currentMethod, currentUri, sw); + + currentUri = ResolveRedirectUri(currentUri, response.Headers.Location); + if (response.StatusCode == HttpStatusCode.SeeOther) + { + currentMethod = "GET"; + body = string.Empty; + } + + continue; + } + + var read = await ReadContentAsync(response.Content, maxResponseBytes, timeoutCts.Token); + if (read.Exceeded) + return Failure($"http response exceeded {maxResponseBytes} bytes", currentMethod, currentUri, sw); + + sw.Stop(); + var metadata = BuildMetadata(response, currentMethod, currentUri, sw.Elapsed.TotalMilliseconds); + return new OutboundHttpResponse + { + Success = response.IsSuccessStatusCode, + Output = read.Text, + Error = response.IsSuccessStatusCode ? string.Empty : BuildHttpErrorMessage(response, read.Text), + Metadata = metadata, + }; + } + catch (OperationCanceledException) when (!ct.IsCancellationRequested) + { + return Failure($"http timeout after {timeoutMs}ms", currentMethod, currentUri, sw); + } + catch (Exception ex) when (!ct.IsCancellationRequested) + { + return Failure(ex.Message, currentMethod, currentUri, sw); + } + } + } + + private static HttpRequestMessage BuildRequestMessage( + string method, + Uri uri, + OutboundHttpRequest request, + string body, + string contentType) + { + var message = new HttpRequestMessage(new HttpMethod(method), uri); + message.Options.Set(AllowPrivateNetworkOption, request.AllowPrivateNetwork); + + foreach (var (key, value) in request.Headers) + { + if (string.IsNullOrWhiteSpace(key)) + continue; + + message.Headers.TryAddWithoutValidation(key.Trim(), value); + } + + if (!string.IsNullOrWhiteSpace(request.Authorization) && + AuthenticationHeaderValue.TryParse(request.Authorization.Trim(), out var authorization)) + { + message.Headers.Authorization = authorization; + } + + if (!string.IsNullOrWhiteSpace(request.IdempotencyKey) && + !message.Headers.Contains("Idempotency-Key")) + { + message.Headers.TryAddWithoutValidation("Idempotency-Key", request.IdempotencyKey.Trim()); + } + + if (!message.Headers.Accept.Any()) + message.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("application/json")); + + if (!string.IsNullOrEmpty(body) || method is not ("GET" or "HEAD")) + message.Content = new StringContent(body, Encoding.UTF8, contentType); + + return message; + } + + private static HttpClient CreateHardenedHttpClient(IOutboundHttpDnsResolver dnsResolver) => + new( + new SocketsHttpHandler + { + AllowAutoRedirect = false, + ConnectCallback = (context, ct) => ConnectValidatedSocketAsync(context, dnsResolver, ct), + }); + + private static async ValueTask ConnectValidatedSocketAsync( + SocketsHttpConnectionContext context, + IOutboundHttpDnsResolver dnsResolver, + CancellationToken ct) + { + var allowPrivateNetwork = + context.InitialRequestMessage?.Options.TryGetValue(AllowPrivateNetworkOption, out var configured) == true && + configured; + var addresses = await ResolveAddressesAsync(context.DnsEndPoint.Host, dnsResolver, ct); + if (addresses.Count == 0) + throw new HttpRequestException($"http DNS resolution returned no addresses for '{context.DnsEndPoint.Host}'"); + + if (!allowPrivateNetwork && addresses.Any(IsBlockedAddress)) + throw new HttpRequestException($"http target '{context.DnsEndPoint.Host}' resolved to a blocked destination"); + + var errors = new List(); + foreach (var address in addresses) + { + var socket = new Socket(address.AddressFamily, SocketType.Stream, ProtocolType.Tcp) + { + NoDelay = true, + }; + try + { + await socket.ConnectAsync(new IPEndPoint(address, context.DnsEndPoint.Port), ct); + return new NetworkStream(socket, ownsSocket: true); + } + catch (Exception ex) when (!ct.IsCancellationRequested) + { + errors.Add(ex); + socket.Dispose(); + } + } + + throw new HttpRequestException( + $"http connection failed for '{context.DnsEndPoint.Host}'", + new AggregateException(errors)); + } + + private async ValueTask ValidateTargetAsync( + Uri uri, + OutboundHttpRequest request, + CancellationToken ct) + { + if (!uri.IsAbsoluteUri) + return "http target must be absolute"; + + if (!string.Equals(uri.Scheme, Uri.UriSchemeHttps, StringComparison.OrdinalIgnoreCase) && + !(request.AllowInsecureHttp && + string.Equals(uri.Scheme, Uri.UriSchemeHttp, StringComparison.OrdinalIgnoreCase))) + { + return "http target must use HTTPS"; + } + + var addresses = await ResolveAddressesAsync(uri.Host, ct); + if (addresses.Count == 0) + return $"http DNS resolution returned no addresses for '{uri.Host}'"; + + if (!request.AllowPrivateNetwork && addresses.Any(IsBlockedAddress)) + return $"http target '{uri.Host}' resolved to a blocked destination"; + + return string.Empty; + } + + private async ValueTask> ResolveAddressesAsync( + string host, + CancellationToken ct) + { + if (IPAddress.TryParse(host, out var literal)) + return [literal]; + + return await ResolveAddressesAsync(host, _dnsResolver, ct); + } + + private static async ValueTask> ResolveAddressesAsync( + string host, + IOutboundHttpDnsResolver dnsResolver, + CancellationToken ct) + { + if (IPAddress.TryParse(host, out var literal)) + return [literal]; + + return await dnsResolver.GetHostAddressesAsync(host, ct); + } + + private static bool IsBlockedAddress(IPAddress address) + { + if (IPAddress.IsLoopback(address) || + IPAddress.IsLoopback(address.MapToIPv6()) || + address.Equals(IPAddress.Any) || + address.Equals(IPAddress.IPv6Any) || + address.Equals(IPAddress.None) || + address.Equals(IPAddress.IPv6None)) + { + return true; + } + + if (address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetwork || + address.IsIPv4MappedToIPv6) + { + var bytes = (address.IsIPv4MappedToIPv6 ? address.MapToIPv4() : address).GetAddressBytes(); + return bytes[0] == 0 || + bytes[0] == 10 || + bytes[0] == 127 || + bytes[0] == 169 && bytes[1] == 254 || + bytes[0] == 172 && bytes[1] is >= 16 and <= 31 || + bytes[0] == 192 && bytes[1] == 168 || + bytes[0] == 100 && bytes[1] is >= 64 and <= 127 || + bytes[0] >= 224; + } + + return address.IsIPv6LinkLocal || + address.IsIPv6Multicast || + address.IsIPv6SiteLocal || + IsUniqueLocalIPv6(address); + } + + private static bool IsUniqueLocalIPv6(IPAddress address) + { + var bytes = address.GetAddressBytes(); + return bytes.Length == 16 && (bytes[0] & 0xfe) == 0xfc; + } + + private static Uri ApplyQuery(Uri uri, IReadOnlyDictionary query) + { + if (query.Count == 0) + return uri; + + var builder = new UriBuilder(uri); + var existing = builder.Query; + if (existing.StartsWith("?", StringComparison.Ordinal)) + existing = existing[1..]; + + var appended = string.Join( + "&", + query + .Where(pair => !string.IsNullOrWhiteSpace(pair.Key)) + .Select(pair => + string.Concat( + Uri.EscapeDataString(pair.Key.Trim()), + "=", + Uri.EscapeDataString(pair.Value ?? string.Empty)))); + builder.Query = string.IsNullOrWhiteSpace(existing) + ? appended + : string.Concat(existing, "&", appended); + return builder.Uri; + } + + private static Uri ResolveRedirectUri(Uri currentUri, Uri location) => + location.IsAbsoluteUri ? location : new Uri(currentUri, location); + + private static bool IsRedirect(HttpStatusCode statusCode) => + statusCode is HttpStatusCode.MovedPermanently or + HttpStatusCode.Found or + HttpStatusCode.SeeOther or + HttpStatusCode.TemporaryRedirect or + HttpStatusCode.PermanentRedirect or + HttpStatusCode.MultipleChoices; + + private static async ValueTask<(bool Exceeded, string Text)> ReadContentAsync( + HttpContent content, + int maxResponseBytes, + CancellationToken ct) + { + await using var stream = await content.ReadAsStreamAsync(ct); + using var buffer = new MemoryStream(Math.Min(maxResponseBytes, 8192)); + var chunk = new byte[8192]; + while (true) + { + var read = await stream.ReadAsync(chunk, ct); + if (read == 0) + break; + + if (buffer.Length + read > maxResponseBytes) + return (true, string.Empty); + + buffer.Write(chunk, 0, read); + } + + return (false, Encoding.UTF8.GetString(buffer.ToArray())); + } + + private static Dictionary BuildMetadata( + HttpResponseMessage response, + string method, + Uri uri, + double durationMs) => + new() + { + ["connector.http.status_code"] = ((int)response.StatusCode).ToString(), + ["connector.http.reason"] = response.ReasonPhrase ?? string.Empty, + ["connector.http.method"] = method, + ["connector.http.url"] = uri.ToString(), + ["connector.http.duration_ms"] = durationMs.ToString("F2"), + }; + + private static OutboundHttpResponse Failure(string error) => + new() + { + Success = false, + Error = error, + }; + + private static OutboundHttpResponse Failure( + string error, + string method, + Uri uri, + Stopwatch sw) + { + sw.Stop(); + return new OutboundHttpResponse + { + Success = false, + Error = error, + Metadata = new Dictionary + { + ["connector.http.method"] = method, + ["connector.http.url"] = uri.ToString(), + ["connector.http.duration_ms"] = sw.Elapsed.TotalMilliseconds.ToString("F2"), + }, + }; + } + + private static string BuildHttpErrorMessage(HttpResponseMessage response, string body) + { + var baseMessage = $"{(int)response.StatusCode} {response.ReasonPhrase}".Trim(); + var detail = TryExtractErrorDetail(body); + return string.IsNullOrWhiteSpace(detail) ? baseMessage : $"{baseMessage}: {detail}"; + } + + private static string TryExtractErrorDetail(string body) + { + if (string.IsNullOrWhiteSpace(body)) + return string.Empty; + + try + { + using var document = JsonDocument.Parse(body); + if (document.RootElement.ValueKind != JsonValueKind.Object) + return RawBodyPreview(body); + + if (document.RootElement.TryGetProperty("description", out var description) && + description.ValueKind == JsonValueKind.String) + { + return description.GetString()?.Trim() ?? string.Empty; + } + + if (document.RootElement.TryGetProperty("error", out var error) && + error.ValueKind == JsonValueKind.String) + { + return error.GetString()?.Trim() ?? string.Empty; + } + } + catch (JsonException) + { + return RawBodyPreview(body); + } + + return RawBodyPreview(body); + } + + private static string RawBodyPreview(string body) + { + var trimmed = body.Trim(); + return trimmed.Length <= 200 ? trimmed : $"{trimmed[..200]}..."; + } + + private static string NormalizeMethod(string? method) => + string.IsNullOrWhiteSpace(method) ? "GET" : method.Trim().ToUpperInvariant(); + + private static int ClampOrDefault(int value, int min, int max, int fallback) + { + if (value <= 0) + return fallback; + return Math.Clamp(value, min, max); + } + + private sealed class DefaultOutboundHttpDnsResolver : IOutboundHttpDnsResolver + { + public async ValueTask> GetHostAddressesAsync( + string host, + CancellationToken ct = default) => + await Dns.GetHostAddressesAsync(host, ct); + } +} diff --git a/src/Aevatar.Studio.Domain/Studio/Compatibility/WorkflowCompatibilityProfile.cs b/src/Aevatar.Studio.Domain/Studio/Compatibility/WorkflowCompatibilityProfile.cs index eec7bfa936..3e3b7b954c 100644 --- a/src/Aevatar.Studio.Domain/Studio/Compatibility/WorkflowCompatibilityProfile.cs +++ b/src/Aevatar.Studio.Domain/Studio/Compatibility/WorkflowCompatibilityProfile.cs @@ -102,7 +102,7 @@ public bool IsSupportedWorkflowCallLifecycle(string? value) } public bool ShouldMirrorTimeoutMsToParameters(string? canonicalType) => - ToCanonicalType(canonicalType) is "wait_signal" or "connector_call" or "llm_call" or "human_input" or "human_approval"; + ToCanonicalType(canonicalType) is "wait_signal" or "http_request" or "connector_call" or "llm_call" or "human_input" or "human_approval"; public string FormatRootFields() => WorkflowYamlRootSchema.FormatAuthorableRootFields(); @@ -163,6 +163,7 @@ private static WorkflowCompatibilityProfile CreateAevatarV1() "workflow_call", "dynamic_workflow", "vote", + "http_request", "connector_call", "emit", "human_input", @@ -222,10 +223,10 @@ private static WorkflowCompatibilityProfile CreateAevatarV1() ["bridge_call"] = "connector_call", ["cli_call"] = "connector_call", ["mcp_call"] = "connector_call", - ["http_get"] = "connector_call", - ["http_post"] = "connector_call", - ["http_put"] = "connector_call", - ["http_delete"] = "connector_call", + ["http_get"] = "http_request", + ["http_post"] = "http_request", + ["http_put"] = "http_request", + ["http_delete"] = "http_request", ["vote_consensus"] = "vote", }), CanonicalStepTypes = canonicalTypes, diff --git a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto index bf137a706f..2806f779ad 100644 --- a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto +++ b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto @@ -485,6 +485,31 @@ message WorkflowStepParameters WorkflowHumanApprovalOptions human_approval = 9; WorkflowExternalApprovalWaitOptions external_approval = 10; WorkflowConnectorApprovalOptions connector_approval = 11; + WorkflowHttpRequestOptions http_request = 12; +} + +message WorkflowHttpRequestAuthentication +{ + string scheme = 1; + string secret_ref = 2; + string header_name = 3; + string header_value_prefix = 4; +} + +message WorkflowHttpRequestOptions +{ + string method = 1; + string url = 2; + map query = 3; + map headers = 4; + string body_mode = 5; + string body = 6; + WorkflowHttpRequestAuthentication authentication = 7; + int32 timeout_ms = 8; + int32 max_response_bytes = 9; + int32 max_redirects = 10; + bool allow_insecure_http = 11; + int32 max_request_bytes = 12; } message StepRequestEvent { diff --git a/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs b/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs index ee9e3184c9..031aeb4707 100644 --- a/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs +++ b/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs @@ -1913,6 +1913,7 @@ private StepRequestEvent BuildStepRequest( ApplyHumanApprovalOptions(request, step.HumanApprovalOptions); ApplyExternalApprovalOptions(request, step.ExternalApprovalOptions, state); ApplyConnectorApprovalOptions(request, step.ConnectorApprovalOptions, state); + ApplyHttpRequestOptions(request, step.HttpRequestOptions, state); ApplyInteractionPresentation(request, step.Presentation, state); return request; @@ -2027,6 +2028,46 @@ private void ApplyConnectorApprovalOptions( }; } + private void ApplyHttpRequestOptions( + StepRequestEvent request, + HttpRequestOptionsDefinition? options, + WorkflowExecutionKernelState state) + { + if (options == null) + return; + + var payload = new WorkflowHttpRequestOptions + { + Method = EvaluateOption(options.Method, state), + Url = EvaluateOption(options.Url, state), + BodyMode = EvaluateOption(options.BodyMode, state), + Body = EvaluateOption(options.Body, state), + TimeoutMs = options.TimeoutMs, + MaxRequestBytes = options.MaxRequestBytes, + MaxResponseBytes = options.MaxResponseBytes, + MaxRedirects = options.MaxRedirects, + AllowInsecureHttp = options.AllowInsecureHttp, + }; + + foreach (var (key, value) in options.Query) + payload.Query[EvaluateOption(key, state)] = EvaluateOption(value, state); + foreach (var (key, value) in options.Headers) + payload.Headers[EvaluateOption(key, state)] = EvaluateOption(value, state); + + if (options.Authentication != null) + { + payload.Authentication = new WorkflowHttpRequestAuthentication + { + Scheme = EvaluateOption(options.Authentication.Scheme, state), + SecretRef = EvaluateOption(options.Authentication.SecretRef, state), + HeaderName = EvaluateOption(options.Authentication.HeaderName, state), + HeaderValuePrefix = EvaluateOption(options.Authentication.HeaderValuePrefix, state), + }; + } + + (request.StepParameters ??= new WorkflowStepParameters()).HttpRequest = payload; + } + private static string NormalizeOptionToken(string? value) => string.IsNullOrWhiteSpace(value) ? string.Empty diff --git a/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.cs b/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.cs index ad622609d2..91fe899364 100644 --- a/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.cs +++ b/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.cs @@ -2,12 +2,15 @@ using Aevatar.Foundation.Abstractions; using Aevatar.Foundation.Core; using Aevatar.Foundation.Abstractions.Connectors; +using Aevatar.Foundation.Abstractions.Credentials; using Aevatar.Foundation.Abstractions.EventModules; using Aevatar.Foundation.Abstractions.Runtime.Callbacks; +using Aevatar.Foundation.Core.Connectors; using Aevatar.Workflow.Core.Execution; using Aevatar.Workflow.Abstractions.Execution; using Aevatar.Workflow.Core.Primitives; using Aevatar.Workflow.Abstractions.Credentials; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace Aevatar.Workflow.Core.Modules; @@ -23,15 +26,21 @@ public sealed partial class ConnectorCallModule : IEventModule "connector_call"; @@ -117,12 +126,20 @@ await HandleApprovalStatusCheckAsync( var request = envelope.Payload.Unpack(); var canonicalStepType = WorkflowPrimitiveCatalog.ToCanonicalType(request.StepType); var isSecureStep = string.Equals(canonicalStepType, "secure_connector_call", StringComparison.OrdinalIgnoreCase); + var isHttpRequestStep = string.Equals(canonicalStepType, "http_request", StringComparison.OrdinalIgnoreCase); if (!string.Equals(canonicalStepType, "connector_call", StringComparison.OrdinalIgnoreCase) && - !isSecureStep) + !isSecureStep && + !isHttpRequestStep) { return; } + if (isHttpRequestStep) + { + await HandleHttpRequestAsync(envelope, request, ctx, ct); + return; + } + var connectorName = WorkflowParameterValueParser.GetString( request.Parameters, string.Empty, @@ -212,6 +229,63 @@ await StartAttemptAsync( ct); } + private async Task HandleHttpRequestAsync( + EventEnvelope envelope, + StepRequestEvent request, + IWorkflowExecutionContext ctx, + CancellationToken ct) + { + var httpRequest = request.StepParameters?.HttpRequest; + if (httpRequest == null) + { + await PublishFailureAsync(ctx, request, "http_request missing typed parameters", ct); + return; + } + + if (ContainsRawAuthorizationHeader(httpRequest.Headers)) + { + await PublishFailureAsync(ctx, request, "http_request authentication must use authentication.secret_ref", ct); + return; + } + + var retry = ParseBoundedInt(request.Parameters.GetValueOrDefault("retry", "0"), 0, 5, 0); + var timeoutMs = ParseBoundedInt( + httpRequest.TimeoutMs > 0 ? httpRequest.TimeoutMs.ToString() : request.Parameters.GetValueOrDefault("timeout_ms", "30000"), + 100, + 300_000, + 30_000); + var attempts = Math.Max(1, retry + 1); + var runId = string.IsNullOrEmpty(request.RunId) + ? envelope.Propagation?.CorrelationId ?? string.Empty + : request.RunId; + var onErrorContinue = string.Equals( + request.Parameters.GetValueOrDefault("on_error", "fail"), + "continue", + StringComparison.OrdinalIgnoreCase); + var normalizedHttpRequest = httpRequest.Clone(); + normalizedHttpRequest.TimeoutMs = timeoutMs; + var connector = new HttpRequestWorkflowConnector( + ResolveOutboundHttpRequestExecutor(ctx), + ResolveCredentialProvider(ctx), + normalizedHttpRequest); + await StartAttemptAsync( + envelope, + request, + runId, + "http_request", + BuildHttpRequestOperation(httpRequest), + connector, + attempt: 1, + attempts, + timeoutMs, + onErrorContinue, + isSecureStep: false, + ctx, + ct, + stepType: "http_request", + httpRequest: normalizedHttpRequest); + } + private async Task HandleTimeoutFiredAsync( WorkflowConnectorTimeoutFiredEvent evt, EventEnvelope envelope, @@ -319,7 +393,12 @@ private async Task HandleAttemptCompletedAsync( "ConnectorCall: step={StepId} connector={Connector} attempt={Attempt}/{Attempts} failed: {Error}", pending.StepId, pending.ConnectorName, pending.Attempt, pending.Attempts, errorText); var nextRequest = BuildRetryStepRequest(pending); - var connector = await _connectorResolver.ResolveAsync(ctx, pending.ConnectorName, ct); + var connector = string.Equals(pending.StepType, "http_request", StringComparison.OrdinalIgnoreCase) + ? new HttpRequestWorkflowConnector( + ResolveOutboundHttpRequestExecutor(ctx), + ResolveCredentialProvider(ctx), + pending.HttpRequest?.Clone() ?? new WorkflowHttpRequestOptions()) + : await _connectorResolver.ResolveAsync(ctx, pending.ConnectorName, ct); if (connector == null) { await PublishPendingCompletionAsync( @@ -347,7 +426,9 @@ await StartAttemptAsync( pending.OnErrorContinue, pending.SecureStep, ctx, - ct); + ct, + stepType: pending.StepType, + httpRequest: pending.HttpRequest); return; } @@ -376,7 +457,9 @@ private async Task StartAttemptAsync( bool isSecureStep, IWorkflowExecutionContext ctx, CancellationToken ct, - string approvalActionId = "") + string approvalActionId = "", + string stepType = "", + WorkflowHttpRequestOptions? httpRequest = null) { var pending = await RegisterPendingAsync( envelope, @@ -392,7 +475,9 @@ private async Task StartAttemptAsync( isSecureStep, ctx, ct, - approvalActionId); + approvalActionId, + stepType, + httpRequest); var requestMetadata = new Dictionary(StringComparer.Ordinal); WorkflowRequestMetadataRuntimeContextAccess.CopyRequestMetadata(ctx, requestMetadata); var connectorRequest = new ConnectorRequest @@ -457,7 +542,9 @@ private async Task RegisterPendingAsync( bool isSecureStep, IWorkflowExecutionContext ctx, CancellationToken ct, - string approvalActionId = "") + string approvalActionId = "", + string stepType = "", + WorkflowHttpRequestOptions? httpRequest = null) { var operationId = BuildOperationId(runId, request.StepId, attempt, request.ExecutionId, ResolveOriginEnvelopeId(envelope)); var callbackId = RuntimeCallbackKeyComposer.BuildCallbackId( @@ -485,6 +572,10 @@ private async Task RegisterPendingAsync( IdempotencyKey = request.IdempotencyKey ?? string.Empty, ApprovalActionId = approvalActionId, RequestDispatched = false, + StepType = string.IsNullOrWhiteSpace(stepType) + ? WorkflowPrimitiveCatalog.ToCanonicalType(request.StepType) + : stepType, + HttpRequest = httpRequest?.Clone(), }; if (string.IsNullOrWhiteSpace(approvalActionId)) { @@ -607,7 +698,9 @@ private static StepRequestEvent BuildRetryStepRequest(PendingConnectorCallState var request = new StepRequestEvent { StepId = pending.StepId, - StepType = pending.SecureStep ? "secure_connector_call" : "connector_call", + StepType = string.IsNullOrWhiteSpace(pending.StepType) + ? (pending.SecureStep ? "secure_connector_call" : "connector_call") + : pending.StepType, RunId = pending.RunId, Input = pending.Input, ExecutionId = pending.ExecutionId, @@ -615,9 +708,196 @@ private static StepRequestEvent BuildRetryStepRequest(PendingConnectorCallState }; foreach (var (key, value) in pending.Parameters) request.Parameters[key] = value; + if (string.Equals(pending.StepType, "http_request", StringComparison.OrdinalIgnoreCase) && + pending.HttpRequest != null) + { + request.StepParameters = new WorkflowStepParameters + { + HttpRequest = pending.HttpRequest.Clone(), + }; + } return request; } + private ICredentialProvider? ResolveCredentialProvider(IWorkflowExecutionContext ctx) => + _credentialProvider ?? ctx.Services.GetService(); + + private IOutboundHttpRequestExecutor ResolveOutboundHttpRequestExecutor(IWorkflowExecutionContext ctx) => + _outboundHttpRequestExecutor ?? + ctx.Services.GetService() ?? + new DefaultOutboundHttpRequestExecutor(); + + private static bool ContainsRawAuthorizationHeader(IEnumerable> headers) => + headers.Any(pair => string.Equals(pair.Key, "Authorization", StringComparison.OrdinalIgnoreCase)); + + private static string BuildHttpRequestOperation(WorkflowHttpRequestOptions options) + { + var method = string.IsNullOrWhiteSpace(options.Method) ? "GET" : options.Method.Trim().ToUpperInvariant(); + if (!Uri.TryCreate(options.Url?.Trim(), UriKind.Absolute, out var uri)) + return method; + + return $"{method} {uri.GetLeftPart(UriPartial.Path)}"; + } + + private sealed class HttpRequestWorkflowConnector( + IOutboundHttpRequestExecutor executor, + ICredentialProvider? credentialProvider, + WorkflowHttpRequestOptions options) : IConnector + { + public string Name => "http_request"; + + public string Type => "http"; + + public async Task ExecuteAsync(ConnectorRequest request, CancellationToken ct = default) + { + var headers = CopyHeadersWithoutContentType(options.Headers, out var contentType); + if (ContainsRawAuthorizationHeader(headers)) + return Failure("http_request authentication must use authentication.secret_ref"); + + var authorization = string.Empty; + var redactions = new List(); + var authResult = await ApplyAuthenticationAsync(headers, redactions, ct); + if (!authResult.Success) + return Failure(authResult.Error); + authorization = authResult.Authorization; + + var response = await executor.ExecuteAsync(new OutboundHttpRequest + { + Method = options.Method, + Url = options.Url, + Query = options.Query.ToDictionary(kv => kv.Key, kv => kv.Value, StringComparer.Ordinal), + Headers = headers, + Authorization = authorization, + IdempotencyKey = request.IdempotencyKey, + Body = ResolveBody(request.Payload), + ContentType = contentType, + TimeoutMs = options.TimeoutMs, + MaxRequestBytes = options.MaxRequestBytes, + MaxResponseBytes = options.MaxResponseBytes, + MaxRedirects = options.MaxRedirects, + AllowInsecureHttp = options.AllowInsecureHttp, + AllowPrivateNetwork = false, + }, ct); + + return new ConnectorResponse + { + Success = response.Success, + Output = RedactSecrets(response.Output, redactions), + Error = RedactSecrets(response.Error, redactions), + Metadata = response.Metadata.ToDictionary( + kv => kv.Key, + kv => RedactSecrets(kv.Value, redactions), + StringComparer.Ordinal), + }; + } + + private async Task<(bool Success, string Authorization, string Error)> ApplyAuthenticationAsync( + IDictionary headers, + List redactions, + CancellationToken ct) + { + var authentication = options.Authentication; + if (authentication == null || + string.IsNullOrWhiteSpace(authentication.Scheme) && + string.IsNullOrWhiteSpace(authentication.SecretRef)) + { + return (true, string.Empty, string.Empty); + } + + var scheme = string.IsNullOrWhiteSpace(authentication.Scheme) + ? "bearer" + : authentication.Scheme.Trim().ToLowerInvariant(); + if (string.Equals(scheme, "none", StringComparison.Ordinal)) + return (true, string.Empty, string.Empty); + + if (string.IsNullOrWhiteSpace(authentication.SecretRef)) + return (false, string.Empty, "http_request authentication requires authentication.secret_ref"); + if (credentialProvider == null) + return (false, string.Empty, "http_request credential provider is unavailable"); + + var secret = await credentialProvider.ResolveAsync(authentication.SecretRef.Trim(), ct); + if (string.IsNullOrEmpty(secret)) + return (false, string.Empty, "http_request authentication secret_ref could not be resolved"); + + redactions.Add(secret); + return scheme switch + { + "bearer" => (true, $"Bearer {secret}", string.Empty), + "header" or "secret_ref_header" => ApplyHeaderAuthentication(headers, authentication, secret), + _ => (false, string.Empty, $"unsupported http_request authentication scheme '{authentication.Scheme}'"), + }; + } + + private static (bool Success, string Authorization, string Error) ApplyHeaderAuthentication( + IDictionary headers, + WorkflowHttpRequestAuthentication authentication, + string secret) + { + if (string.IsNullOrWhiteSpace(authentication.HeaderName)) + return (false, string.Empty, "http_request header authentication requires authentication.header_name"); + + headers[authentication.HeaderName.Trim()] = string.Concat( + authentication.HeaderValuePrefix ?? string.Empty, + secret); + return (true, string.Empty, string.Empty); + } + + private string ResolveBody(string inheritedPayload) + { + var mode = options.BodyMode?.Trim(); + if (string.Equals(mode, "none", StringComparison.OrdinalIgnoreCase)) + return string.Empty; + if (string.Equals(mode, "input", StringComparison.OrdinalIgnoreCase) || + string.Equals(mode, "inherit", StringComparison.OrdinalIgnoreCase)) + { + return inheritedPayload ?? string.Empty; + } + + return options.Body ?? string.Empty; + } + + private static Dictionary CopyHeadersWithoutContentType( + IEnumerable> source, + out string contentType) + { + contentType = string.Empty; + var headers = new Dictionary(StringComparer.OrdinalIgnoreCase); + foreach (var (key, value) in source) + { + if (string.IsNullOrWhiteSpace(key)) + continue; + + if (string.Equals(key, "Content-Type", StringComparison.OrdinalIgnoreCase)) + { + contentType = value ?? string.Empty; + continue; + } + + headers[key.Trim()] = value ?? string.Empty; + } + + return headers; + } + + private static ConnectorResponse Failure(string error) => + new() + { + Success = false, + Error = error, + }; + + private static string RedactSecrets(string value, IReadOnlyCollection secrets) + { + if (string.IsNullOrEmpty(value) || secrets.Count == 0) + return value ?? string.Empty; + + var redacted = value; + foreach (var secret in secrets.Where(secret => !string.IsNullOrEmpty(secret))) + redacted = redacted.Replace(secret, "[redacted]", StringComparison.Ordinal); + return redacted; + } + } + private static double ParseDuration(WorkflowConnectorAttemptCompletedEvent evt) { return double.TryParse( diff --git a/src/workflow/Aevatar.Workflow.Core/Primitives/StepDefinition.cs b/src/workflow/Aevatar.Workflow.Core/Primitives/StepDefinition.cs index 19e5aeade7..7046319af2 100644 --- a/src/workflow/Aevatar.Workflow.Core/Primitives/StepDefinition.cs +++ b/src/workflow/Aevatar.Workflow.Core/Primitives/StepDefinition.cs @@ -43,6 +43,8 @@ public sealed class StepDefinition public ConnectorApprovalOptionsDefinition? ConnectorApprovalOptions { get; init; } + public HttpRequestOptionsDefinition? HttpRequestOptions { get; init; } + /// /// 下一步骤 ID,用于线性流程控制。 /// @@ -138,6 +140,44 @@ public sealed class ConnectorApprovalOptionsDefinition public string? PolicyReason { get; init; } } +public sealed class HttpRequestOptionsDefinition +{ + public string? Method { get; init; } + + public string? Url { get; init; } + + public Dictionary Query { get; init; } = new(StringComparer.Ordinal); + + public Dictionary Headers { get; init; } = new(StringComparer.OrdinalIgnoreCase); + + public string? BodyMode { get; init; } + + public string? Body { get; init; } + + public HttpRequestAuthenticationDefinition? Authentication { get; init; } + + public int TimeoutMs { get; init; } + + public int MaxRequestBytes { get; init; } + + public int MaxResponseBytes { get; init; } + + public int MaxRedirects { get; init; } + + public bool AllowInsecureHttp { get; init; } +} + +public sealed class HttpRequestAuthenticationDefinition +{ + public string? Scheme { get; init; } + + public string? SecretRef { get; init; } + + public string? HeaderName { get; init; } + + public string? HeaderValuePrefix { get; init; } +} + public sealed class StepRetryPolicy { /// 最大尝试次数(含首次)。最小 1,最大 10。 diff --git a/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowParser.cs b/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowParser.cs index 1e1ec8673f..2f1f8c3484 100644 --- a/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowParser.cs +++ b/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowParser.cs @@ -195,6 +195,7 @@ private static StepDefinition MapStep(RawStep s) HumanApprovalOptions = MapHumanApprovalOptions(canonicalType, parameters), ExternalApprovalOptions = MapExternalApprovalOptions(canonicalType, parameters), ConnectorApprovalOptions = MapConnectorApprovalOptions(canonicalType, parameters), + HttpRequestOptions = MapHttpRequestOptions(canonicalType, parameters), Next = s.Next, Compensation = NormalizeText(s.Compensation), Children = s.Children?.Select(MapStep).ToList(), @@ -952,7 +953,7 @@ s.TimeoutMs is not null && } private static bool ShouldLiftTimeoutMsToParameter(string canonicalType) => - canonicalType is "wait_signal" or "connector_call" or "secure_connector_call" or "llm_call" or "human_input" or "secure_input" or "human_approval"; + canonicalType is "wait_signal" or "http_request" or "connector_call" or "secure_connector_call" or "llm_call" or "human_input" or "secure_input" or "human_approval"; private static void LiftTransformOperationParameters( string canonicalType, @@ -1095,6 +1096,156 @@ private static void LiftTransformOperationParameters( }; } + private static HttpRequestOptionsDefinition? MapHttpRequestOptions( + string canonicalType, + IReadOnlyDictionary parameters) + { + if (!string.Equals(canonicalType, "http_request", StringComparison.Ordinal)) + return null; + + var authentication = MapHttpRequestAuthentication(parameters); + return new HttpRequestOptionsDefinition + { + Method = GetParameter(parameters, "method").Trim(), + Url = GetParameter(parameters, "url").Trim(), + Query = ReadStringMapParameter(parameters, "query", "query.", StringComparer.Ordinal), + Headers = ReadStringMapParameter(parameters, "headers", "header.", StringComparer.OrdinalIgnoreCase), + BodyMode = GetParameter(parameters, "body_mode", "stdin_mode").Trim(), + Body = GetParameter(parameters, "body", "payload", "stdin_value"), + Authentication = authentication, + TimeoutMs = ParseNonNegativeInt(GetParameter(parameters, "timeout_ms")), + MaxRequestBytes = ParseNonNegativeInt(GetParameter(parameters, "max_request_bytes")), + MaxResponseBytes = ParseNonNegativeInt(GetParameter(parameters, "max_response_bytes")), + MaxRedirects = ParseNonNegativeInt(GetParameter(parameters, "max_redirects")), + AllowInsecureHttp = ParseBool(GetParameter(parameters, "allow_insecure_http")), + }; + } + + private static HttpRequestAuthenticationDefinition? MapHttpRequestAuthentication( + IReadOnlyDictionary parameters) + { + var auth = ReadObjectParameter(parameters, "authentication", "auth"); + var scheme = ReadObjectString(auth, "scheme"); + var secretRef = ReadObjectString(auth, "secret_ref", "secretRef", "credential_ref", "credentialRef"); + var headerName = ReadObjectString(auth, "header_name", "headerName"); + var headerValuePrefix = ReadObjectString(auth, "header_value_prefix", "headerValuePrefix"); + + scheme = FirstNonWhiteSpace(scheme, GetParameter(parameters, "authentication.scheme", "auth.scheme", "auth_scheme")); + secretRef = FirstNonWhiteSpace( + secretRef, + GetParameter( + parameters, + "authentication.secret_ref", + "auth.secret_ref", + "credential_ref", + "secret_ref")); + headerName = FirstNonWhiteSpace( + headerName, + GetParameter(parameters, "authentication.header_name", "auth.header_name", "auth_header_name")); + headerValuePrefix = FirstNonWhiteSpace( + headerValuePrefix, + GetParameter( + parameters, + "authentication.header_value_prefix", + "auth.header_value_prefix", + "auth_header_value_prefix")); + + if (string.IsNullOrWhiteSpace(scheme) && + string.IsNullOrWhiteSpace(secretRef) && + string.IsNullOrWhiteSpace(headerName) && + string.IsNullOrWhiteSpace(headerValuePrefix)) + { + return null; + } + + return new HttpRequestAuthenticationDefinition + { + Scheme = scheme.Trim(), + SecretRef = secretRef.Trim(), + HeaderName = headerName.Trim(), + HeaderValuePrefix = headerValuePrefix, + }; + } + + private static Dictionary ReadStringMapParameter( + IReadOnlyDictionary parameters, + string mapKey, + string prefix, + StringComparer comparer) + { + var values = new Dictionary(comparer); + var map = ReadObjectParameter(parameters, mapKey); + foreach (var (key, value) in map) + { + if (!string.IsNullOrWhiteSpace(key)) + values[key.Trim()] = value; + } + + foreach (var (key, value) in parameters) + { + if (!key.StartsWith(prefix, StringComparison.OrdinalIgnoreCase)) + continue; + + var mapEntryKey = key[prefix.Length..].Trim(); + if (!string.IsNullOrWhiteSpace(mapEntryKey)) + values[mapEntryKey] = value; + } + + return values; + } + + private static Dictionary ReadObjectParameter( + IReadOnlyDictionary parameters, + params string[] keys) + { + var raw = GetParameter(parameters, keys); + if (string.IsNullOrWhiteSpace(raw)) + return []; + + try + { + using var document = JsonDocument.Parse(raw); + if (document.RootElement.ValueKind != JsonValueKind.Object) + return []; + + var values = new Dictionary(StringComparer.Ordinal); + foreach (var property in document.RootElement.EnumerateObject()) + values[property.Name] = JsonElementToString(property.Value); + return values; + } + catch (JsonException) + { + return []; + } + } + + private static string ReadObjectString( + IReadOnlyDictionary values, + params string[] keys) => + GetParameter(values, keys).Trim(); + + private static string FirstNonWhiteSpace(params string[] values) + { + foreach (var value in values) + { + if (!string.IsNullOrWhiteSpace(value)) + return value; + } + + return string.Empty; + } + + private static string JsonElementToString(JsonElement element) => + element.ValueKind switch + { + JsonValueKind.String => element.GetString() ?? string.Empty, + JsonValueKind.Number => element.GetRawText(), + JsonValueKind.True => "true", + JsonValueKind.False => "false", + JsonValueKind.Null => string.Empty, + _ => element.GetRawText(), + }; + private static int ParseNonNegativeInt(string? value) => int.TryParse(value, NumberStyles.Integer, CultureInfo.InvariantCulture, out var parsed) && parsed > 0 ? parsed diff --git a/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowPrimitiveCatalog.cs b/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowPrimitiveCatalog.cs index d70d725113..96dd7c1334 100644 --- a/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowPrimitiveCatalog.cs +++ b/src/workflow/Aevatar.Workflow.Core/Primitives/WorkflowPrimitiveCatalog.cs @@ -29,10 +29,10 @@ public static class WorkflowPrimitiveCatalog ["bridge_call"] = "connector_call", ["cli_call"] = "connector_call", ["mcp_call"] = "connector_call", - ["http_get"] = "connector_call", - ["http_post"] = "connector_call", - ["http_put"] = "connector_call", - ["http_delete"] = "connector_call", + ["http_get"] = "http_request", + ["http_post"] = "http_request", + ["http_put"] = "http_request", + ["http_delete"] = "http_request", ["secure_connector"] = "secure_connector_call", ["secret_input"] = "secure_input", ["vote_consensus"] = "vote", @@ -49,7 +49,7 @@ public static class WorkflowPrimitiveCatalog private static readonly string[] CapabilityPrimitives = [ - "llm_call", "tool_call", "connector_call", "secure_connector_call", + "llm_call", "tool_call", "http_request", "connector_call", "secure_connector_call", "evaluate", "reflect", "human_input", "secure_input", "human_approval", "wait_signal", "emit", "parallel", "race", "map_reduce", "vote", "foreach", "dynamic_workflow", "lease", @@ -114,6 +114,7 @@ public static bool IsSideEffectingPrimitive(string stepType) var canonical = ToCanonicalType(stepType); // Saga v1.1 provisions only the current external-dispatch primitives; future primitives can opt in explicitly. return string.Equals(canonical, "tool_call", StringComparison.Ordinal) || + string.Equals(canonical, "http_request", StringComparison.Ordinal) || string.Equals(canonical, "connector_call", StringComparison.Ordinal) || string.Equals(canonical, "secure_connector_call", StringComparison.Ordinal); } diff --git a/src/workflow/Aevatar.Workflow.Core/ServiceCollectionExtensions.cs b/src/workflow/Aevatar.Workflow.Core/ServiceCollectionExtensions.cs index f1c806ac2e..9309145919 100644 --- a/src/workflow/Aevatar.Workflow.Core/ServiceCollectionExtensions.cs +++ b/src/workflow/Aevatar.Workflow.Core/ServiceCollectionExtensions.cs @@ -2,6 +2,7 @@ using Aevatar.Foundation.Abstractions.Connectors; using Aevatar.Foundation.Abstractions.EventModules; using Aevatar.Foundation.Abstractions.EventSourcing; +using Aevatar.Foundation.Core.Connectors; using Aevatar.Foundation.Core.TypeSystem; using Aevatar.Workflow.Abstractions.Execution; using Aevatar.Workflow.Core.Execution; @@ -28,6 +29,7 @@ public static IServiceCollection AddAevatarWorkflow(this IServiceCollection serv services.TryAddSingleton, WorkflowModuleFactory>(); services.TryAddSingleton(); services.TryAddSingleton(); + services.TryAddSingleton(_ => new DefaultOutboundHttpRequestExecutor()); services.TryAddSingleton(); services.TryAddEnumerable(ServiceDescriptor.Singleton< ICommittedStatePublicationHook, diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowCoreModulePack.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowCoreModulePack.cs index 3964314e0d..3c0bb8efc7 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowCoreModulePack.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowCoreModulePack.cs @@ -20,7 +20,7 @@ public sealed class WorkflowCoreModulePack : IWorkflowModulePack WorkflowModuleRegistration.Create("map_reduce", "mapreduce"), WorkflowModuleRegistration.Create("llm_call"), WorkflowModuleRegistration.Create("tool_call"), - WorkflowModuleRegistration.Create("connector_call", "bridge_call", "secure_connector_call", "secure_connector"), + WorkflowModuleRegistration.Create("connector_call", "bridge_call", "secure_connector_call", "secure_connector", "http_request"), WorkflowModuleRegistration.Create("transform"), WorkflowModuleRegistration.Create("retrieve_facts"), WorkflowModuleRegistration.Create("wait_signal", "wait"), diff --git a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto index e11a1c0658..481f6a16a9 100644 --- a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto +++ b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto @@ -623,6 +623,8 @@ message PendingConnectorCallState string idempotency_key = 17; string approval_action_id = 18; bool request_dispatched = 19; + string step_type = 20; + aevatar.workflow.WorkflowHttpRequestOptions http_request = 21; } message ConnectorCallProtectedParameter diff --git a/src/workflow/Aevatar.Workflow.Infrastructure/Capabilities/WorkflowInfrastructureCapabilitiesProvider.cs b/src/workflow/Aevatar.Workflow.Infrastructure/Capabilities/WorkflowInfrastructureCapabilitiesProvider.cs index 196169c427..360930e5eb 100644 --- a/src/workflow/Aevatar.Workflow.Infrastructure/Capabilities/WorkflowInfrastructureCapabilitiesProvider.cs +++ b/src/workflow/Aevatar.Workflow.Infrastructure/Capabilities/WorkflowInfrastructureCapabilitiesProvider.cs @@ -32,6 +32,23 @@ internal sealed class WorkflowInfrastructureCapabilitiesProvider : IWorkflowCapa ["secure_connector_call"] = new( "Invokes an external connector with secure payload handling.", ConnectorParameterDescriptors()), + ["http_request"] = new( + "Sends a direct outbound HTTP request with typed URL, method, body, limits, and secret-backed authentication.", + [ + new PrimitiveParameterDescriptor("method", "string", true, "HTTP method.", DefaultValue: "GET", EnumValuesInput: ["GET", "POST", "PUT", "PATCH", "DELETE"]), + new PrimitiveParameterDescriptor("url", "string", true, "Absolute HTTPS URL."), + new PrimitiveParameterDescriptor("query", "object", false, "Query parameters appended to the URL."), + new PrimitiveParameterDescriptor("headers", "object", false, "Non-authorization request headers."), + new PrimitiveParameterDescriptor("body_mode", "string", false, "Body source: none, raw, or input.", DefaultValue: "none", EnumValuesInput: ["none", "raw", "input"]), + new PrimitiveParameterDescriptor("body", "string", false, "Raw request body when body_mode is raw."), + new PrimitiveParameterDescriptor("authentication", "object", false, "Authentication object containing scheme and secret_ref."), + new PrimitiveParameterDescriptor("timeout_ms", "int", false, "Request timeout in milliseconds.", DefaultValue: "30000"), + new PrimitiveParameterDescriptor("max_request_bytes", "int", false, "Maximum request body bytes.", DefaultValue: "65536"), + new PrimitiveParameterDescriptor("max_response_bytes", "int", false, "Maximum response body bytes.", DefaultValue: "65536"), + new PrimitiveParameterDescriptor("max_redirects", "int", false, "Maximum validated redirects.", DefaultValue: "3"), + new PrimitiveParameterDescriptor("retry", "int", false, "Retry count for failed attempts.", DefaultValue: "0"), + new PrimitiveParameterDescriptor("on_error", "string", false, "Failure behavior.", DefaultValue: "fail", EnumValuesInput: ["fail", "continue"]), + ]), ["llm_call"] = new( "Runs an LLM role step and returns generated output.", [ @@ -252,7 +269,7 @@ private static string InferPrimitiveCategory(string canonicalType) => "guard" or "conditional" or "switch" or "while" or "delay" or "wait_signal" or "checkpoint" or "workflow_loop" or "workflow_yaml_validate" => "control", "foreach" or "parallel" or "race" or "map_reduce" or "workflow_call" or "vote" or "dynamic_workflow" or "self_reschedule" => "composition", "llm_call" or "tool_call" or "evaluate" or "reflect" => "ai", - "connector_call" or "secure_connector_call" or "emit" => "integration", + "connector_call" or "secure_connector_call" or "http_request" or "emit" => "integration", "human_input" or "human_approval" or "secure_input" => "human", _ => "general", }; diff --git a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowCatalogCurrentStateProjector.cs b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowCatalogCurrentStateProjector.cs index 4fa747f2fe..5b417abb47 100644 --- a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowCatalogCurrentStateProjector.cs +++ b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowCatalogCurrentStateProjector.cs @@ -21,6 +21,7 @@ public sealed class WorkflowCatalogCurrentStateProjector "wait_signal", "connector_call", "secure_connector_call", + "http_request", }; private readonly IProjectionWriteDispatcher _writeDispatcher; diff --git a/test/Aevatar.Foundation.Core.Tests/OutboundHttpRequestExecutorTests.cs b/test/Aevatar.Foundation.Core.Tests/OutboundHttpRequestExecutorTests.cs new file mode 100644 index 0000000000..cf8f6b7cee --- /dev/null +++ b/test/Aevatar.Foundation.Core.Tests/OutboundHttpRequestExecutorTests.cs @@ -0,0 +1,225 @@ +using System.Net; +using System.Reflection; +using System.Text; +using Aevatar.Foundation.Abstractions.Connectors; +using Aevatar.Foundation.Core.Connectors; +using FluentAssertions; + +namespace Aevatar.Foundation.Core.Tests; + +public sealed class OutboundHttpRequestExecutorTests +{ + [Fact] + public async Task ExecuteAsync_ShouldRejectPrivateIpLiteralBeforeSending() + { + var handler = new RecordingHttpMessageHandler(_ => + new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("{}"), + }); + var executor = new DefaultOutboundHttpRequestExecutor( + new HttpClient(handler), + new StaticDnsResolver(IPAddress.Parse("93.184.216.34"))); + + var response = await executor.ExecuteAsync(new OutboundHttpRequest + { + Method = "GET", + Url = "https://127.0.0.1/admin", + TimeoutMs = 5000, + }); + + response.Success.Should().BeFalse(); + response.Error.Should().Contain("blocked destination"); + handler.Requests.Should().BeEmpty(); + } + + [Fact] + public async Task ExecuteAsync_ShouldValidateRedirectTarget() + { + var handler = new RecordingHttpMessageHandler(request => + { + if (request.RequestUri!.AbsoluteUri == "https://api.example.com/start") + { + return new HttpResponseMessage(HttpStatusCode.Redirect) + { + Headers = { Location = new Uri("https://169.254.169.254/latest/meta-data") }, + }; + } + + return new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("""{"ok":true}"""), + }; + }); + var resolver = new StaticDnsResolver( + ("api.example.com", IPAddress.Parse("93.184.216.34")), + ("169.254.169.254", IPAddress.Parse("169.254.169.254"))); + var executor = new DefaultOutboundHttpRequestExecutor(new HttpClient(handler), resolver); + + var response = await executor.ExecuteAsync(new OutboundHttpRequest + { + Method = "GET", + Url = "https://api.example.com/start", + TimeoutMs = 5000, + MaxRedirects = 3, + }); + + response.Success.Should().BeFalse(); + response.Error.Should().Contain("blocked destination"); + handler.Requests.Should().ContainSingle() + .Which.RequestUri!.AbsoluteUri.Should().Be("https://api.example.com/start"); + } + + [Fact] + public void DefaultExecutor_ShouldValidateDnsAtSocketConnectionBoundary() + { + var executor = new DefaultOutboundHttpRequestExecutor(); + + var client = typeof(DefaultOutboundHttpRequestExecutor) + .GetField("_client", BindingFlags.Instance | BindingFlags.NonPublic)! + .GetValue(executor) + .Should() + .BeOfType() + .Subject; + var handler = typeof(HttpMessageInvoker) + .GetField("_handler", BindingFlags.Instance | BindingFlags.NonPublic)! + .GetValue(client) + .Should() + .BeOfType() + .Subject; + + handler.ConnectCallback.Should().NotBeNull(); + } + + [Fact] + public async Task ExecuteAsync_ShouldFailWhenResponseExceedsLimit() + { + var handler = new RecordingHttpMessageHandler(_ => + new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("abcdef", Encoding.UTF8, "text/plain"), + }); + var executor = new DefaultOutboundHttpRequestExecutor( + new HttpClient(handler), + new StaticDnsResolver(IPAddress.Parse("93.184.216.34"))); + + var response = await executor.ExecuteAsync(new OutboundHttpRequest + { + Method = "GET", + Url = "https://api.example.com/data", + TimeoutMs = 5000, + MaxResponseBytes = 5, + }); + + response.Success.Should().BeFalse(); + response.Error.Should().Contain("response exceeded 5 bytes"); + response.Output.Should().BeEmpty(); + } + + [Fact] + public async Task ExecuteAsync_ShouldRejectRequestBodyThatExceedsLimitBeforeSending() + { + var handler = new RecordingHttpMessageHandler(_ => + new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("{}"), + }); + var executor = new DefaultOutboundHttpRequestExecutor( + new HttpClient(handler), + new StaticDnsResolver(IPAddress.Parse("93.184.216.34"))); + var request = new OutboundHttpRequest + { + Method = "POST", + Url = "https://api.example.com/data", + Body = "abcdef", + ContentType = "text/plain", + TimeoutMs = 5000, + MaxRequestBytes = 5, + }; + + var response = await executor.ExecuteAsync(request); + + response.Success.Should().BeFalse(); + response.Error.Should().Contain("request exceeded 5 bytes"); + handler.Requests.Should().BeEmpty(); + } + + [Fact] + public async Task ExecuteAsync_ShouldReturnStatusAndBodyForNonSuccess() + { + var handler = new RecordingHttpMessageHandler(_ => + new HttpResponseMessage(HttpStatusCode.ServiceUnavailable) + { + Content = new StringContent("down", Encoding.UTF8, "text/plain"), + ReasonPhrase = "Service Unavailable", + }); + var executor = new DefaultOutboundHttpRequestExecutor( + new HttpClient(handler), + new StaticDnsResolver(IPAddress.Parse("93.184.216.34"))); + + var response = await executor.ExecuteAsync(new OutboundHttpRequest + { + Method = "POST", + Url = "https://api.example.com/data", + Body = "{}", + ContentType = "application/json", + TimeoutMs = 5000, + }); + + response.Success.Should().BeFalse(); + response.Error.Should().Contain("503 Service Unavailable"); + response.Output.Should().Be("down"); + response.Metadata["connector.http.status_code"].Should().Be("503"); + handler.Requests.Should().ContainSingle(); + } + + private sealed class RecordingHttpMessageHandler( + Func responseFactory) : HttpMessageHandler + { + public List Requests { get; } = []; + + protected override Task SendAsync( + HttpRequestMessage request, + CancellationToken cancellationToken) + { + _ = cancellationToken; + Requests.Add(request); + return Task.FromResult(responseFactory(request)); + } + } + + private sealed class StaticDnsResolver : IOutboundHttpDnsResolver + { + private readonly Dictionary _addresses; + + public StaticDnsResolver(params IPAddress[] addresses) + { + _addresses = new Dictionary(StringComparer.OrdinalIgnoreCase) + { + ["*"] = addresses, + }; + } + + public StaticDnsResolver(params (string Host, IPAddress Address)[] addresses) + { + _addresses = addresses + .GroupBy(entry => entry.Host, StringComparer.OrdinalIgnoreCase) + .ToDictionary( + group => group.Key, + group => group.Select(entry => entry.Address).ToArray(), + StringComparer.OrdinalIgnoreCase); + } + + public ValueTask> GetHostAddressesAsync( + string host, + CancellationToken ct = default) + { + _ = ct; + return ValueTask.FromResult>( + _addresses.TryGetValue(host, out var addresses) || + _addresses.TryGetValue("*", out addresses) + ? addresses + : []); + } + } +} diff --git a/test/Aevatar.Integration.Tests/ConnectorCallModuleHttpRequestTests.cs b/test/Aevatar.Integration.Tests/ConnectorCallModuleHttpRequestTests.cs new file mode 100644 index 0000000000..078c896768 --- /dev/null +++ b/test/Aevatar.Integration.Tests/ConnectorCallModuleHttpRequestTests.cs @@ -0,0 +1,285 @@ +using Aevatar.Foundation.Abstractions; +using Aevatar.Foundation.Abstractions.Connectors; +using Aevatar.Foundation.Abstractions.Credentials; +using Aevatar.Foundation.Core; +using Aevatar.Workflow.Abstractions.Execution; +using Aevatar.Workflow.Core; +using Aevatar.Workflow.Core.Connectors; +using Aevatar.Workflow.Core.Modules; +using FluentAssertions; +using Google.Protobuf; +using Google.Protobuf.WellKnownTypes; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Aevatar.Integration.Tests; + +[Trait("Category", "Integration")] +[Trait("Feature", "ConnectorCallModuleHttpRequest")] +public sealed class ConnectorCallModuleHttpRequestTests +{ + [Fact] + public async Task HandleAsync_WhenHttpRequestUsesTypedOptions_ShouldExecuteWithoutConnectorRegistry() + { + var executor = new RecordingOutboundHttpRequestExecutor(new OutboundHttpResponse + { + Success = true, + Output = """{"ok":true}""", + Metadata = new Dictionary + { + ["connector.http.status_code"] = "200", + ["connector.http.method"] = "GET", + ["connector.http.url"] = "https://api.example.com/q1000", + }, + }); + var module = new ConnectorCallModule(new RegistryBackedWorkflowConnectorResolver(new ConfiguredConnectorRegistry())); + var ctx = CreateContext(executor, new StubCredentialProvider(("scope-secret:q1000-token", "sk-http-secret"))); + var request = new StepRequestEvent + { + StepId = "s-http", + RunId = "run-http", + StepType = "http_request", + ExecutionId = "exec-http", + IdempotencyKey = "idem-http", + StepParameters = new WorkflowStepParameters + { + HttpRequest = new WorkflowHttpRequestOptions + { + Method = "GET", + Url = "https://api.example.com/q1000", + TimeoutMs = 20_000, + MaxRequestBytes = 4096, + MaxResponseBytes = 65_536, + MaxRedirects = 2, + Authentication = new WorkflowHttpRequestAuthentication + { + Scheme = "bearer", + SecretRef = "scope-secret:q1000-token", + }, + }, + }, + }; + request.StepParameters.HttpRequest.Headers["X-Trace"] = "trace-123"; + + await HandleAndDrainAsync(module, Envelope(request), ctx); + + executor.Requests.Should().ContainSingle(); + var outbound = executor.Requests.Single(); + outbound.Method.Should().Be("GET"); + outbound.Url.Should().Be("https://api.example.com/q1000"); + outbound.Authorization.Should().Be("Bearer sk-http-secret"); + outbound.Headers.Should().Contain("X-Trace", "trace-123"); + outbound.Headers.Should().NotContainKey("Authorization"); + outbound.IdempotencyKey.Should().Be("idem-http"); + outbound.MaxRequestBytes.Should().Be(4096); + + var completed = ctx.Published.Should().ContainSingle().Subject.evt.Should().BeOfType().Subject; + completed.Success.Should().BeTrue(); + completed.Output.Should().Be("""{"ok":true}"""); + completed.Error.Should().BeEmpty(); + completed.Annotations["connector.name"].Should().Be("http_request"); + completed.Annotations["connector.type"].Should().Be("http"); + completed.Annotations["connector.http.status_code"].Should().Be("200"); + completed.Annotations.Values.Should().NotContain(value => value.Contains("sk-http-secret", StringComparison.Ordinal)); + + var state = ctx.LoadState("connector_call"); + state.PendingByOperationId.Should().BeEmpty(); + } + + [Fact] + public async Task HandleAsync_WhenHttpRequestUsesRawAuthorizationHeader_ShouldFailBeforeExecution() + { + var executor = new RecordingOutboundHttpRequestExecutor(new OutboundHttpResponse + { + Success = true, + Output = """{"ok":true}""", + }); + var module = new ConnectorCallModule(new RegistryBackedWorkflowConnectorResolver(new ConfiguredConnectorRegistry())); + var ctx = CreateContext(executor, new StubCredentialProvider()); + var request = new StepRequestEvent + { + StepId = "s-raw-auth", + RunId = "run-raw-auth", + StepType = "http_request", + StepParameters = new WorkflowStepParameters + { + HttpRequest = new WorkflowHttpRequestOptions + { + Method = "GET", + Url = "https://api.example.com/q1000", + }, + }, + }; + request.StepParameters.HttpRequest.Headers["Authorization"] = "Bearer raw-token"; + + await module.HandleAsync(Envelope(request), ctx, CancellationToken.None); + + executor.Requests.Should().BeEmpty(); + var completed = ctx.Published.Should().ContainSingle().Subject.evt.Should().BeOfType().Subject; + completed.Success.Should().BeFalse(); + completed.Error.Should().Be("http_request authentication must use authentication.secret_ref"); + completed.Error.Should().NotContain("raw-token"); + } + + [Fact] + public async Task HandleAsync_WhenHttpRequestRetries_ShouldPreservePrimitiveIdentityAndTypedOptions() + { + var executor = new RecordingOutboundHttpRequestExecutor( + new OutboundHttpResponse + { + Success = false, + Error = "503 Service Unavailable", + }, + new OutboundHttpResponse + { + Success = true, + Output = """{"ok":true}""", + }); + var module = new ConnectorCallModule(new RegistryBackedWorkflowConnectorResolver(new ConfiguredConnectorRegistry())); + var ctx = CreateContext(executor, new StubCredentialProvider(("retry-secret", "retry-token"))); + var request = new StepRequestEvent + { + StepId = "s-http-retry", + RunId = "run-http-retry", + StepType = "http_request", + Input = "payload", + StepParameters = new WorkflowStepParameters + { + HttpRequest = new WorkflowHttpRequestOptions + { + Method = "POST", + Url = "https://api.example.com/retry", + Body = """{"source":"test"}""", + BodyMode = "raw", + TimeoutMs = 10_000, + MaxRequestBytes = 2048, + MaxResponseBytes = 1024, + MaxRedirects = 1, + Authentication = new WorkflowHttpRequestAuthentication + { + Scheme = "bearer", + SecretRef = "retry-secret", + }, + }, + }, + }; + request.Parameters["retry"] = "1"; + + await module.HandleAsync(Envelope(request), ctx, CancellationToken.None); + + var pending = ctx.LoadState("connector_call") + .PendingByOperationId + .Values + .Should() + .ContainSingle() + .Subject; + pending.StepType.Should().Be("http_request"); + pending.HttpRequest.Should().NotBeNull(); + pending.HttpRequest.Url.Should().Be("https://api.example.com/retry"); + + await DrainConnectorContinuationsAsync(module, ctx); + + executor.Requests.Should().HaveCount(2); + executor.Requests.Select(sent => sent.Url).Should().OnlyContain(url => url == "https://api.example.com/retry"); + executor.Requests.Select(sent => sent.Authorization).Should().OnlyContain(value => value == "Bearer retry-token"); + executor.Requests.Select(sent => sent.MaxRequestBytes).Should().OnlyContain(value => value == 2048); + + var completed = ctx.Published.Should().ContainSingle().Subject.evt.Should().BeOfType().Subject; + completed.Success.Should().BeTrue(); + completed.Output.Should().Be("""{"ok":true}"""); + completed.Annotations["connector.name"].Should().Be("http_request"); + completed.Annotations["connector.attempts"].Should().Be("2"); + completed.Annotations.Values.Should().NotContain(value => value.Contains("retry-token", StringComparison.Ordinal)); + } + + private static TestEventHandlerContext CreateContext( + IOutboundHttpRequestExecutor outboundHttpRequestExecutor, + ICredentialProvider credentialProvider) + { + var services = new ServiceCollection() + .AddSingleton(outboundHttpRequestExecutor) + .AddSingleton(credentialProvider) + .BuildServiceProvider(); + return new TestEventHandlerContext( + services, + new TestAgent("connector-http-request-test-agent"), + NullLogger.Instance); + } + + private static EventEnvelope Envelope(IMessage evt, string? correlationId = null) => + new() + { + Id = Guid.NewGuid().ToString("N"), + Timestamp = Timestamp.FromDateTime(DateTime.UtcNow), + Payload = Any.Pack(evt), + Route = EnvelopeRouteSemantics.CreateTopologyPublication("test-publisher", TopologyAudience.Self), + Propagation = new EnvelopePropagation + { + CorrelationId = correlationId ?? string.Empty, + }, + }; + + private static async Task HandleAndDrainAsync( + ConnectorCallModule module, + EventEnvelope envelope, + TestEventHandlerContext ctx) + { + await module.HandleAsync(envelope, ctx, CancellationToken.None); + await DrainConnectorContinuationsAsync(module, ctx); + } + + private static async Task DrainConnectorContinuationsAsync( + ConnectorCallModule module, + TestEventHandlerContext ctx) + { + for (var index = 0; index < ctx.Published.Count; index++) + { + if (ctx.Published[index].evt is not WorkflowConnectorAttemptCompletedEvent completed) + continue; + + ctx.Published.RemoveAt(index); + index--; + await module.HandleAsync(Envelope(completed), ctx, CancellationToken.None); + } + } + + private sealed class RecordingOutboundHttpRequestExecutor(params OutboundHttpResponse[] responses) + : IOutboundHttpRequestExecutor + { + private readonly Queue _responses = new(responses); + + public List Requests { get; } = []; + + public Task ExecuteAsync( + OutboundHttpRequest request, + CancellationToken ct = default) + { + _ = ct; + Requests.Add(request); + return Task.FromResult(_responses.Count > 0 + ? _responses.Dequeue() + : new OutboundHttpResponse + { + Success = true, + Output = "{}", + }); + } + } + + private sealed class StubCredentialProvider(params (string Ref, string Secret)[] credentials) : ICredentialProvider + { + private readonly Dictionary _credentials = credentials.ToDictionary( + entry => entry.Ref, + entry => entry.Secret, + StringComparer.Ordinal); + + public Task ResolveAsync(string credentialRef, CancellationToken ct = default) + { + _ = ct; + return Task.FromResult( + _credentials.TryGetValue(credentialRef, out var secret) + ? secret + : null); + } + } +} diff --git a/test/Aevatar.Studio.Tests/WorkflowCompatibilityProfileTests.cs b/test/Aevatar.Studio.Tests/WorkflowCompatibilityProfileTests.cs index 343778d258..63a9faa451 100644 --- a/test/Aevatar.Studio.Tests/WorkflowCompatibilityProfileTests.cs +++ b/test/Aevatar.Studio.Tests/WorkflowCompatibilityProfileTests.cs @@ -141,8 +141,8 @@ private static string BuildYamlWithRootField(string rootField) => [InlineData("loop", "while")] [InlineData("sub_workflow", "workflow_call")] [InlineData("foreach_llm", "foreach")] - [InlineData("http_get", "connector_call")] - [InlineData("http_post", "connector_call")] + [InlineData("http_get", "http_request")] + [InlineData("http_post", "http_request")] [InlineData("mcp_call", "connector_call")] [InlineData("sleep", "delay")] [InlineData("publish", "emit")] @@ -156,6 +156,7 @@ public void ToCanonicalType_ShouldResolveAliases(string alias, string expected) [InlineData("transform")] [InlineData("conditional")] [InlineData("llm_call")] + [InlineData("http_request")] [InlineData("workflow_call")] public void ToCanonicalType_ShouldReturnCanonicalAsIs(string type) { @@ -182,6 +183,7 @@ public void ToCanonicalType_ShouldBeCaseInsensitiveAndTrim(string value, string [Theory] [InlineData("transform", true)] [InlineData("loop", true)] + [InlineData("http_request", true)] [InlineData("actor_send", true)] [InlineData("workflow_loop", true)] [InlineData("nonexistent", false)] @@ -241,6 +243,7 @@ public void IsSupportedWorkflowCallLifecycle_ShouldValidateLifecycles(string? va [Theory] [InlineData("wait_signal", true)] [InlineData("connector_call", true)] + [InlineData("http_request", true)] [InlineData("llm_call", true)] [InlineData("human_input", true)] [InlineData("human_approval", true)] diff --git a/test/Aevatar.Studio.Tests/WorkflowDocumentNormalizerTests.cs b/test/Aevatar.Studio.Tests/WorkflowDocumentNormalizerTests.cs index 047bacba3d..7d0369935c 100644 --- a/test/Aevatar.Studio.Tests/WorkflowDocumentNormalizerTests.cs +++ b/test/Aevatar.Studio.Tests/WorkflowDocumentNormalizerTests.cs @@ -171,6 +171,7 @@ public void NormalizeForExport_ShouldApplyHttpGetDefaults() Steps = [new StepModel { Id = "s1", Type = "http_get" }], }; var result = _normalizer.NormalizeForExport(doc); + result.Steps[0].Type.Should().Be("http_request"); result.Steps[0].Parameters.Should().ContainKey("method"); result.Steps[0].Parameters["method"]!.ToWorkflowScalarString().Should().Be("GET"); } diff --git a/test/Aevatar.Workflow.Core.Tests/Primitives/WorkflowParserConfigurationTests.cs b/test/Aevatar.Workflow.Core.Tests/Primitives/WorkflowParserConfigurationTests.cs index bb9b9de226..409d13010a 100644 --- a/test/Aevatar.Workflow.Core.Tests/Primitives/WorkflowParserConfigurationTests.cs +++ b/test/Aevatar.Workflow.Core.Tests/Primitives/WorkflowParserConfigurationTests.cs @@ -218,19 +218,19 @@ public void Parse_WhenUsingErgonomicAliases_ShouldCanonicalizeAndApplyDefaults() - id: s_http_get type: http_get parameters: - connector: demo_http + url: https://api.example.com/get - id: s_http_post type: http_post parameters: - connector: demo_http + url: https://api.example.com/post - id: s_http_put type: http_put parameters: - connector: demo_http + url: https://api.example.com/put - id: s_http_delete type: http_delete parameters: - connector: demo_http + url: https://api.example.com/delete - id: s_cli type: cli_call parameters: @@ -252,14 +252,23 @@ public void Parse_WhenUsingErgonomicAliases_ShouldCanonicalizeAndApplyDefaults() var workflow = new WorkflowParser().Parse(yaml); - workflow.Steps.First(s => s.Id == "s_http_get").Type.Should().Be("connector_call"); - workflow.Steps.First(s => s.Id == "s_http_get").Parameters["method"].Should().Be("GET"); - workflow.Steps.First(s => s.Id == "s_http_post").Type.Should().Be("connector_call"); - workflow.Steps.First(s => s.Id == "s_http_post").Parameters["method"].Should().Be("POST"); - workflow.Steps.First(s => s.Id == "s_http_put").Type.Should().Be("connector_call"); - workflow.Steps.First(s => s.Id == "s_http_put").Parameters["method"].Should().Be("PUT"); - workflow.Steps.First(s => s.Id == "s_http_delete").Type.Should().Be("connector_call"); - workflow.Steps.First(s => s.Id == "s_http_delete").Parameters["method"].Should().Be("DELETE"); + var httpGet = workflow.Steps.First(s => s.Id == "s_http_get"); + httpGet.Type.Should().Be("http_request"); + httpGet.Parameters["method"].Should().Be("GET"); + httpGet.HttpRequestOptions.Should().NotBeNull(); + httpGet.HttpRequestOptions!.Url.Should().Be("https://api.example.com/get"); + + var httpPost = workflow.Steps.First(s => s.Id == "s_http_post"); + httpPost.Type.Should().Be("http_request"); + httpPost.Parameters["method"].Should().Be("POST"); + + var httpPut = workflow.Steps.First(s => s.Id == "s_http_put"); + httpPut.Type.Should().Be("http_request"); + httpPut.Parameters["method"].Should().Be("PUT"); + + var httpDelete = workflow.Steps.First(s => s.Id == "s_http_delete"); + httpDelete.Type.Should().Be("http_request"); + httpDelete.Parameters["method"].Should().Be("DELETE"); workflow.Steps.First(s => s.Id == "s_cli").Type.Should().Be("connector_call"); diff --git a/test/Aevatar.Workflow.Core.Tests/WorkflowHttpRequestPrimitiveTests.cs b/test/Aevatar.Workflow.Core.Tests/WorkflowHttpRequestPrimitiveTests.cs new file mode 100644 index 0000000000..9bad8902d0 --- /dev/null +++ b/test/Aevatar.Workflow.Core.Tests/WorkflowHttpRequestPrimitiveTests.cs @@ -0,0 +1,107 @@ +using Aevatar.Workflow.Abstractions; +using Aevatar.Workflow.Core; +using Aevatar.Workflow.Core.Primitives; +using FluentAssertions; + +namespace Aevatar.Workflow.Core.Tests; + +public sealed class WorkflowHttpRequestPrimitiveTests +{ + [Fact] + public void Parse_ShouldMapHttpRequestParametersToTypedOptions() + { + var workflow = new WorkflowParser().Parse(""" + name: direct-http + steps: + - id: fetch_snapshot + type: http_request + parameters: + method: GET + url: "https://api.example.com/q1000" + query: + source: q1000 + headers: + X-Trace: "${input}" + authentication: + scheme: bearer + secret_ref: scope-secret:q1000-token + timeout_ms: "20000" + max_request_bytes: "4096" + max_response_bytes: "65536" + max_redirects: "2" + """); + + var step = workflow.Steps.Should().ContainSingle().Subject; + step.Type.Should().Be("http_request"); + step.Parameters.Should().NotContainKey("connector"); + step.HttpRequestOptions.Should().NotBeNull(); + step.HttpRequestOptions!.Method.Should().Be("GET"); + step.HttpRequestOptions.Url.Should().Be("https://api.example.com/q1000"); + step.HttpRequestOptions.Query.Should().ContainSingle().Which.Should().Be(new KeyValuePair("source", "q1000")); + step.HttpRequestOptions.Headers.Should().ContainSingle().Which.Should().Be(new KeyValuePair("X-Trace", "${input}")); + step.HttpRequestOptions.Authentication.Should().NotBeNull(); + step.HttpRequestOptions.Authentication!.Scheme.Should().Be("bearer"); + step.HttpRequestOptions.Authentication.SecretRef.Should().Be("scope-secret:q1000-token"); + step.HttpRequestOptions.TimeoutMs.Should().Be(20000); + step.HttpRequestOptions.MaxRequestBytes.Should().Be(4096); + step.HttpRequestOptions.MaxResponseBytes.Should().Be(65536); + step.HttpRequestOptions.MaxRedirects.Should().Be(2); + } + + [Theory] + [InlineData("http_get", "GET")] + [InlineData("http_post", "POST")] + [InlineData("http_put", "PUT")] + [InlineData("http_delete", "DELETE")] + public void Parse_ShouldMapHttpAliasesToHttpRequest(string alias, string method) + { + var workflow = new WorkflowParser().Parse($$""" + name: direct-http-alias + steps: + - id: fetch + type: {{alias}} + parameters: + url: "https://api.example.com/items" + """); + + var step = workflow.Steps.Should().ContainSingle().Subject; + step.Type.Should().Be("http_request"); + step.HttpRequestOptions.Should().NotBeNull(); + step.HttpRequestOptions!.Method.Should().Be(method); + step.HttpRequestOptions.Url.Should().Be("https://api.example.com/items"); + } + + [Fact] + public void HttpRequest_ShouldBeSideEffectingWithoutConnectorDependency() + { + WorkflowPrimitiveCatalog.ToCanonicalType("http_request").Should().Be("http_request"); + WorkflowPrimitiveCatalog.IsSideEffectingPrimitive("http_request").Should().BeTrue(); + + var yaml = """ + name: direct-http-dependencies + steps: + - id: fetch + type: http_request + parameters: + method: GET + url: "https://api.example.com/items" + """; + + var result = new WorkflowGAgent().EvaluateAuthorizationDependencies(yaml); + + result.Should().NotBeNull(); + result!.ConnectorCapabilityRefs.Should().BeEmpty(); + result.ServiceGrantPolicy.Should().Be(WorkflowServiceGrantPolicy.NotRequiredNoExternalService); + } + + [Fact] + public void WorkflowStepParameters_ShouldExposeTypedHttpRequestOptions() + { + WorkflowStepParameters.Descriptor.Fields.InDeclarationOrder() + .Should() + .Contain(field => field.Name == "http_request" && field.FieldNumber == 12); + WorkflowHttpRequestOptions.Descriptor.Fields.InDeclarationOrder() + .Should() + .Contain(field => field.Name == "max_request_bytes"); + } +} diff --git a/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs index 1af0cd97ed..d1f45d058b 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs @@ -599,7 +599,7 @@ public async Task HandleChat_ShouldUseDelegationCredentialAndAuthenticatedScope( { capturedCommand = command; return Task.FromResult( - CommandInteractionResult + CommandInteractionResult .Failure(WorkflowChatRunStartError.WorkflowBindingMismatch)); }, }; diff --git a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowInfrastructureCoverageTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowInfrastructureCoverageTests.cs index 336e7f703e..9351b0cec5 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowInfrastructureCoverageTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowInfrastructureCoverageTests.cs @@ -1223,6 +1223,7 @@ [new CustomModulePack()], workflow.AuthorityStateVersion == 7); var primitiveNames = document.Primitives.Select(primitive => primitive.Name).ToList(); primitiveNames.Should().Contain("connector_call"); + primitiveNames.Should().Contain("http_request"); primitiveNames.Should().Contain("llm_call"); primitiveNames.Should().Contain("tool_call"); document.Primitives.Should().Contain(primitive =>