外部 API

使用外部 API 执行已部署的工作流、查询执行日志并接收完成 Webhook。

认证

工作流执行和 Logs API 请求需要在 X-API-Key 请求头中传入 API 密钥:

curl -H "x-api-key: YOUR_API_KEY" \
  https://www.tradinggoose.ai/api/v1/logs?workspaceId=YOUR_WORKSPACE_ID

你可以在 TradingGoose 仪表盘的用户设置中生成 API 密钥。

执行工作流

向已部署的工作流发送 POST 请求。input 对象会传递到其 API 触发器。

curl -X POST \
  https://www.tradinggoose.ai/api/workflows/WORKFLOW_ID/execute \
  -H "Content-Type: application/json" \
  -H "X-API-Key: YOUR_API_KEY" \
  -d '{
    "input": {
      "topic": "Semiconductor earnings"
    }
  }'

请求体接受以下内容:

  • input — 包含工作流输入的对象。
  • stream — 设为 true 以接收服务器发送事件,而不是单个 JSON 响应。
  • selectedOutputs — 由 blockName.path 字符串组成的数组,例如 "Summarize.content",用于选择流式块输出。块名称必须解析到恰好一个已部署块。

内部执行控制和草稿状态覆盖会被拒绝。该端点始终运行已部署的工作流状态。

如果没有 Response 块,非流式请求会返回公开执行结果:

{
  "success": true,
  "output": {
    "content": "..."
  },
  "metadata": {
    "duration": 842,
    "startTime": "2026-09-03T18:30:00.000Z",
    "endTime": "2026-09-03T18:30:00.842Z"
  }
}

如果工作流包含 Response 块,成功执行时会改用该块配置的 JSON 正文、HTTP 状态和响应头。非流式请求为排队执行最多等待 25 秒;超时返回 HTTP 504。

使用 stream: true 时,响应的内容类型为 text/event-stream,并带有 X-Execution-Id 响应头。SSE 事件的数据为 JSON,但结束标记 [DONE] 是字面文本,不是 JSON:

data: {"blockId":"summarize","chunk":"Partial output"}

data: {"event":"final","data":{"success":true,"output":{"content":"..."},"metadata":{"duration":842}}}

data: [DONE]

调用 JSON.parse 前,先检查 data 字段是否等于 [DONE];如果是,就停止读取,不要将该标记作为 JSON 解析。

错误使用相同的信封结构,其中包含 "event":"error" 和一条 error 消息。

Logs API

每个 Logs API 响应都包含有关你的工作流执行限制和用量的信息:

{
  "limits": {
    "executionRateLimit": {
      "sync": {
        "limit": 60,        // Max sync workflow executions per minute
        "remaining": 58,    // Remaining sync workflow executions
        "resetAt": "..."    // When the window resets
      },
      "async": {
        "limit": 60,        // Max async workflow executions per minute
        "remaining": 59,    // Remaining async workflow executions
        "resetAt": "..."    // When the window resets
      }
    },
    "usage": {
      "currentPeriodCost": 1.234,  // Current billing period usage in USD
      "limit": 10,                  // Usage limit in USD
      "tier": {                     // Current billing tier summary
        "id": "tier_pro",
        "displayName": "Pro"
      },
      "isExceeded": false           // Whether limit is exceeded
    }
  }
}

注意: 响应正文中的速率限制针对工作流执行。调用此 API 端点的速率限制位于响应头中(X-RateLimit-*)。

查询日志

使用丰富的筛选选项查询工作流执行日志。

GET /api/v1/logs

必需参数:

  • workspaceId - 你的工作区 ID

可选筛选条件:

  • workflowIds - 以逗号分隔的工作流 ID
  • folderIds - 以逗号分隔的文件夹 ID
  • triggers - 以逗号分隔的触发器类型:api、webhook、schedule、manual、chat
  • level - 按级别筛选:info、error
  • startDate - 日期范围开始的 ISO 时间戳
  • endDate - 日期范围结束的 ISO 时间戳
  • executionId - 精确匹配执行 ID
  • minDurationMs - 最短执行时长(毫秒)
  • maxDurationMs - 最长执行时长(毫秒)
  • minCost - 最低执行成本
  • maxCost - 最高执行成本
  • model - 按使用的 AI 模型筛选
  • monitorId - 精确匹配监控 ID
  • listing - 经 JSON 编码的规范化标识,包含 listing_id、base_id、quote_id 和 listing_type(default、crypto 或 currency)
  • indicatorId - 精确匹配指标 ID
  • providerId - 精确匹配市场或经纪商提供商 ID
  • interval - 精确匹配监控间隔
  • triggerSource - 以逗号分隔的监控触发器 ID:indicator_trigger、portfolio_state_trigger

分页:

  • limit - 每页结果数(默认:100)
  • cursor - 下一页游标
  • order - 排序方式:desc、asc(默认:desc)

详细级别:

  • details - 响应详细级别:basic、full(默认:basic)
  • includeTraceSpans - 当 details=full 时包含 trace span(默认:false)
  • includeFinalOutput - 当 details=full 时包含最终输出(默认:false)

在 details=full 时,只有启用上述任一包含选项,才会通过 errorMessage 返回失败详情。

{
  "data": [
    {
      "id": "log_abc123",
      "workflowId": "wf_xyz789",
      "executionId": "exec_def456",
      "level": "info",
      "trigger": "api",
      "startedAt": "2025-01-01T12:34:56.789Z",
      "endedAt": "2025-01-01T12:34:57.123Z",
      "totalDurationMs": 334,
      "cost": {
        "total": 0.00234
      },
      "files": null
    }
  ],
  "nextCursor": "eyJzIjoiMjAyNS0wMS0wMVQxMjozNDo1Ni43ODlaIiwiaWQiOiJsb2dfYWJjMTIzIn0",
  "limits": {
    "executionRateLimit": {
      "sync": {
        "limit": 60,
        "remaining": 58,
        "resetAt": "2025-01-01T12:35:56.789Z"
      },
      "async": {
        "limit": 60,
        "remaining": 59,
        "resetAt": "2025-01-01T12:35:56.789Z"
      }
    },
    "usage": {
      "currentPeriodCost": 1.234,
      "limit": 10,
      "tier": {
        "id": "tier_pro",
        "displayName": "Pro"
      },
      "isExceeded": false
    }
  }
}

获取日志详情

获取特定日志条目的详细信息。

GET /api/v1/logs/{id}
{
  "data": {
    "id": "log_abc123",
    "workflowId": "wf_xyz789",
    "executionId": "exec_def456",
    "level": "info",
    "trigger": "api",
    "startedAt": "2025-01-01T12:34:56.789Z",
    "endedAt": "2025-01-01T12:34:57.123Z",
    "totalDurationMs": 334,
    "workflow": {
      "id": "wf_xyz789",
      "name": "My Workflow",
      "description": "Process customer data",
      "color": "#2dd4bf",
      "folderId": "folder_123",
      "userId": "user_123",
      "workspaceId": "workspace_123",
      "createdAt": "2025-01-01T10:00:00.000Z",
      "updatedAt": "2025-01-01T11:00:00.000Z"
    },
    "executionData": {
      "traceSpans": [...],
      "finalOutput": {...}
    },
    "cost": {
      "total": 0.00234,
      "tokens": {
        "prompt": 123,
        "completion": 456,
        "total": 579
      },
      "models": {
        "gpt-4o": {
          "input": 0.001,
          "output": 0.00134,
          "total": 0.00234,
          "tokens": {
            "prompt": 123,
            "completion": 456,
            "total": 579
          }
        }
      }
    },
    "createdAt": "2025-01-01T12:34:57.200Z"
  },
  "limits": {
      "executionRateLimit": {
        "sync": {
          "limit": 60,
          "remaining": 58,
          "resetAt": "2025-01-01T12:35:56.789Z"
        },
        "async": {
          "limit": 60,
          "remaining": 59,
          "resetAt": "2025-01-01T12:35:56.789Z"
        }
      },
      "usage": {
        "currentPeriodCost": 1.234,
        "limit": 10,
        "tier": {
          "id": "tier_pro",
          "displayName": "Pro"
        },
        "isExceeded": false
      }
    }
}

获取执行详情

获取执行详情,包括工作流状态快照。

GET /api/v1/logs/executions/{executionId}
{
  "executionId": "exec_def456",
  "workflowId": "wf_xyz789",
  "workflowState": {
    "blocks": {...},
    "edges": [...],
    "loops": {...},
    "parallels": {...}
  },
  "executionMetadata": {
    "trigger": "api",
    "startedAt": "2025-01-01T12:34:56.789Z",
    "endedAt": "2025-01-01T12:34:57.123Z",
    "totalDurationMs": 334,
    "cost": {...}
  },
  "limits": {
    "executionRateLimit": {
      "sync": {
        "limit": 60,
        "remaining": 58,
        "resetAt": "2025-01-01T12:35:56.789Z"
      },
      "async": {
        "limit": 60,
        "remaining": 59,
        "resetAt": "2025-01-01T12:35:56.789Z"
      }
    },
    "usage": {
      "currentPeriodCost": 1.234,
      "limit": 10,
      "tier": {
        "id": "tier_pro",
        "displayName": "Pro"
      },
      "isExceeded": false
    }
  }
}

Webhook 订阅

工作流执行完成时,获取实时通知。

订阅选项

Webhook 订阅管理不属于 API 密钥日志 API 的一部分。应用程序使用会话认证的工作流路由来管理订阅。一旦工作流具有活动订阅,以下选项控制投递:

可用的配置选项:

  • url: 您的 Webhook 端点 URL
  • secret: 用于 HMAC 签名验证的可选密钥
  • includeFinalOutput: 在负载中包含工作流的最终输出
  • includeTraceSpans: 包含详细的执行跟踪跨度
  • includeRateLimits: 包含工作流所有者的速率限制信息
  • includeUsageData: 包含工作流所有者的使用量和账单数据
  • levelFilter: 要接收的日志级别数组 (info, error)
  • triggerFilter: 要接收的触发器类型数组 (api, webhook, schedule, manual, chat)
  • active: 新建订阅是否立即激活(默认:true)

Webhook 负载

当工作流执行完成时,TradingGoose 会向您的 Webhook URL 发送 POST 请求:

{
  "id": "evt_123",
  "type": "workflow.execution.completed",
  "timestamp": 1735925767890,
  "data": {
    "workflowId": "wf_xyz789",
    "executionId": "exec_def456",
    "status": "success",
    "level": "info",
    "trigger": "api",
    "startedAt": "2025-01-01T12:34:56.789Z",
    "endedAt": "2025-01-01T12:34:57.123Z",
    "totalDurationMs": 334,
    "cost": {
      "total": 0.00234,
      "tokens": {
        "prompt": 123,
        "completion": 456,
        "total": 579
      },
      "models": {
        "gpt-4o": {
          "input": 0.001,
          "output": 0.00134,
          "total": 0.00234,
          "tokens": {
            "prompt": 123,
            "completion": 456,
            "total": 579
          }
        }
      }
    },
    "files": null,
    "errorMessage": "...", // 仅在执行失败且 includeFinalOutput=true 或 includeTraceSpans=true 时包含
    "finalOutput": {...},  // Only if includeFinalOutput=true
    "traceSpans": [...],   // Only if includeTraceSpans=true
    "rateLimits": {...},   // Only if includeRateLimits=true
    "usage": {...}         // Only if includeUsageData=true
  },
  "links": {
    "log": "/v1/logs/log_abc123",
    "execution": "/v1/logs/executions/exec_def456"
  }
}

Webhook 标头

每个 webhook 请求都包含以下请求头:

  • tradinggoose-event:事件类型(始终为 workflow.execution.completed)
  • tradinggoose-timestamp:Unix 时间戳(毫秒)
  • tradinggoose-delivery-id:用于幂等的唯一投递 ID
  • tradinggoose-signature:用于验证的 HMAC-SHA256 签名(若已配置密钥)
  • Idempotency-Key:与投递 ID 相同,用于重复检测

签名验证与重复投递处理

配置 webhook 密钥,并在记录或处理投递前验证签名。重试会保留相同的 Idempotency-Key 和 tradinggoose-delivery-id,但会生成新的负载 id、时间戳和签名。不要使用 event.id 或请求体哈希去重。对于已处理完成的投递,返回 200,不要重复执行其副作用。

这些示例使用持久化 SQLite 数据库,通过原子插入来认领每次投递。将数据库副作用与接收记录放在同一事务中:处理失败时两者一并回滚,重试即可再次处理。唯一约束还涵盖已签名的事件类型、工作流 ID 和执行 ID,因为投递请求头不在签名范围内。因此,对每个接收方而言,一次工作流完成只作为一个业务事件处理;需要独立处理的订阅应使用独立的数据库或消费者作用域。

将标记的注释替换为使用该事务 db 的同步数据库写入;不要启动未等待完成的任务。

Node.js 示例要求 Node.js 24+,并使用内置的 node:sqlite;Python 使用内置的 sqlite3 模块。分别安装 Express 或 Flask,并设置 WEBHOOK_SECRET。将 WEBHOOK_DB_PATH 指向持久化数据库文件(默认:webhooks.sqlite)。请在任何 JSON 中间件之前挂载此路由,以确保签名验证获取原始请求体。

接收方的所有工作进程必须使用同一个持久化接收记录存储。对于不同主机上的副本,请使用共享事务数据库,而不是各自独立的 SQLite 文件。请在重试及重放保留期内保留接收记录。支付或发送邮件等外部副作用无法通过此数据库事务回滚:在确认投递之前,应使用下游服务的幂等机制或事务性发件箱。

import crypto from 'crypto';
import { DatabaseSync } from 'node:sqlite';
import express from 'express';

const app = express();
const db = new DatabaseSync(process.env.WEBHOOK_DB_PATH || 'webhooks.sqlite', { timeout: 5000 });
db.exec(`
  CREATE TABLE IF NOT EXISTS webhook_receipts (
    delivery_id TEXT PRIMARY KEY,
    event_type TEXT NOT NULL,
    workflow_id TEXT NOT NULL,
    execution_id TEXT NOT NULL,
    UNIQUE (event_type, workflow_id, execution_id)
  )
`);
const claimDelivery = db.prepare(`
  INSERT INTO webhook_receipts (delivery_id, event_type, workflow_id, execution_id)
  VALUES (?, ?, ?, ?) ON CONFLICT DO NOTHING
`);

function verifyWebhookSignature(rawBody, signature, secret) {
  if (!signature || !secret) return false;
  const [timestampPart, signaturePart] = signature.split(',');
  if (!timestampPart || !signaturePart) return false;
  const timestamp = timestampPart.replace('t=', '');
  const expectedSignature = signaturePart.replace('v1=', '');

  const computedSignature = crypto
    .createHmac('sha256', secret)
    .update(`${timestamp}.${rawBody}`)
    .digest('hex');

  const computed = Buffer.from(computedSignature, 'hex');
  const expected = Buffer.from(expectedSignature, 'hex');
  return computed.length === expected.length && crypto.timingSafeEqual(computed, expected);
}

app.post('/webhook', express.raw({ type: 'application/json' }), (req, res) => {
  const signature = req.headers['tradinggoose-signature'];
  const rawBody = req.body.toString('utf8');

  if (!verifyWebhookSignature(rawBody, signature, process.env.WEBHOOK_SECRET)) {
    return res.status(401).send('Invalid signature');
  }

  const deliveryHeader = req.headers['tradinggoose-delivery-id'];
  const deliveryId = req.headers['idempotency-key'] || deliveryHeader;
  if (typeof deliveryId !== 'string' || !deliveryId.trim() ||
      (deliveryHeader && deliveryHeader !== deliveryId)) {
    return res.status(400).send('Missing or inconsistent delivery ID');
  }

  let event;
  try {
    event = JSON.parse(rawBody);
  } catch {
    return res.status(400).send('Invalid JSON');
  }
  if (event?.type !== 'workflow.execution.completed' ||
      typeof event.data?.workflowId !== 'string' || !event.data.workflowId ||
      typeof event.data?.executionId !== 'string' || !event.data.executionId) {
    return res.status(400).send('Invalid event');
  }

  try {
    db.exec('BEGIN IMMEDIATE');
    const receipt = claimDelivery.run(deliveryId, event.type, event.data.workflowId, event.data.executionId);
    if (receipt.changes === 0) {
      db.exec('COMMIT');
      return res.sendStatus(200);
    }
    // Apply your database side effects here using db, inside this transaction.
    db.exec('COMMIT');
  } catch (error) {
    if (db.isTransaction) db.exec('ROLLBACK');
    console.error('Webhook processing failed', error);
    return res.sendStatus(500);
  }
  res.sendStatus(200);
});
import hmac
import hashlib
import os
import sqlite3
from contextlib import closing
from flask import Flask, request

app = Flask(__name__)
database_path = os.environ.get('WEBHOOK_DB_PATH', 'webhooks.sqlite')
with closing(sqlite3.connect(database_path)) as db, db:
    db.execute('''
        CREATE TABLE IF NOT EXISTS webhook_receipts (
            delivery_id TEXT PRIMARY KEY,
            event_type TEXT NOT NULL,
            workflow_id TEXT NOT NULL,
            execution_id TEXT NOT NULL,
            UNIQUE (event_type, workflow_id, execution_id)
        )
    ''')

def verify_webhook_signature(raw_body: str, signature: str, secret: str) -> bool:
    if not signature or not secret:
        return False
    timestamp_part, separator, signature_part = signature.partition(',')
    if not separator:
        return False
    timestamp = timestamp_part.replace('t=', '')
    expected_signature = signature_part.replace('v1=', '')

    signature_base = f"{timestamp}.{raw_body}"
    computed_signature = hmac.new(
        secret.encode(),
        signature_base.encode(),
        hashlib.sha256
    ).hexdigest()

    return hmac.compare_digest(computed_signature, expected_signature)

@app.route('/webhook', methods=['POST'])
def webhook():
    signature = request.headers.get('tradinggoose-signature')
    raw_body = request.get_data(as_text=True)

    if not verify_webhook_signature(raw_body, signature, os.environ['WEBHOOK_SECRET']):
        return 'Invalid signature', 401

    event = request.get_json()
    delivery_header = request.headers.get('tradinggoose-delivery-id')
    delivery_id = request.headers.get('Idempotency-Key') or delivery_header
    if not delivery_id or not delivery_id.strip() or (
        delivery_header and delivery_header != delivery_id
    ):
        return 'Missing or inconsistent delivery ID', 400
    data = event.get('data') if isinstance(event, dict) else None
    if not isinstance(data, dict) or event.get('type') != 'workflow.execution.completed' or any(
        not isinstance(data.get(key), str) or not data[key]
        for key in ('workflowId', 'executionId')
    ):
        return 'Invalid event', 400

    try:
        with closing(sqlite3.connect(database_path, timeout=5)) as db, db:
            db.execute('BEGIN IMMEDIATE')
            receipt = db.execute('''
                INSERT INTO webhook_receipts (delivery_id, event_type, workflow_id, execution_id)
                VALUES (?, ?, ?, ?) ON CONFLICT DO NOTHING
            ''', (delivery_id, event['type'], data['workflowId'], data['executionId']))
            if receipt.rowcount == 0:
                return '', 200
            # Apply your database side effects here using db, inside this transaction.
    except Exception:
        app.logger.exception('Webhook processing failed')
        return '', 500
    return '', 200

重试策略

可重试的 webhook 投递最多尝试五次:首次请求加上四次带退避和抖动的重试。

  • 重试延迟:5 秒、15 秒、1 分钟和 3 分钟
  • 抖动:最多额外 10% 的延迟,以防止惊群效应
  • HTTP 5xx、HTTP 429、网络错误和超时会触发重试
  • 投递在 30 秒后超时

webhook 投递为异步处理,不会影响工作流的执行性能。

最佳实践

  1. 轮询策略:轮询日志时,使用基于游标的分页,配合 order=asc 和 startDate 高效获取新日志。

  2. webhook 安全:始终配置 webhook 密钥并验证签名,以确保请求来自 TradingGoose。

  3. 幂等性:在同一个持久化事务中原子记录投递 ID 并执行数据库副作用。对已完成的重复投递返回 200;处理失败时回滚,以便重试能够成功。对于外部副作用,使用下游幂等机制或事务性发件箱。

  4. 隐私:默认情况下,响应中会排除 finalOutput 和 traceSpans。仅在你需要这些数据并了解其隐私影响时才启用。

  5. 速率限制:收到 429 响应时,实施指数退避。查看 Retry-After 请求头以获取建议的等待时间。

速率限制

API 会应用调用方当前计费层级所配置的 API 端点限制。这与日志响应中包含的同步和异步工作流执行限制是分开的。

速率限制信息包含在响应头中:

  • X-RateLimit-Limit:每个窗口的最大请求数
  • X-RateLimit-Remaining:当前窗口内剩余请求数
  • X-RateLimit-Reset:窗口重置时的 ISO 时间戳