Documentation
¶
Index ¶
- func StartWorker(ctx context.Context, svc *Service)
- type BatchTaskCallbackSubscriptionsRequest
- type BatchTaskCallbackSubscriptionsResponse
- type CreateTaskCallbackRequest
- type DeliveryListItem
- type Handler
- func (h *Handler) BatchManageTaskCallbacks(c echo.Context) error
- func (h *Handler) CreateTaskCallback(c echo.Context) error
- func (h *Handler) DeleteTaskCallback(c echo.Context) error
- func (h *Handler) ListManagedTaskCallbacks(c echo.Context) error
- func (h *Handler) ListTaskCallbackDeliveries(c echo.Context) error
- func (h *Handler) ListTaskCallbacks(c echo.Context) error
- func (h *Handler) PauseTaskCallback(c echo.Context) error
- func (h *Handler) RegisterProtected(api *echo.Group, jwtMiddleware echo.MiddlewareFunc)
- func (h *Handler) ResumeTaskCallback(c echo.Context) error
- type Service
- func (s *Service) AttemptAgentWebhookEffect(ctx context.Context, effect db.RunEffectOutbox) runtimepkg.RunEffectAttemptResult
- func (s *Service) AttemptDelivery(ctx context.Context, deliveryID uuid.UUID) error
- func (s *Service) AttemptTaskCallbackDelivery(ctx context.Context, deliveryID uuid.UUID) error
- func (s *Service) AttemptTaskCallbackEffect(ctx context.Context, effect db.RunEffectOutbox) runtimepkg.RunEffectAttemptResult
- func (s *Service) BatchManageTaskCallbackSubscriptions(ctx context.Context, userID uuid.UUID, ...) (*BatchTaskCallbackSubscriptionsResponse, error)
- func (s *Service) ClearWebhook(ctx context.Context, agentID, userID uuid.UUID) error
- func (s *Service) CreateTaskCallbackSubscription(ctx context.Context, runID, userID uuid.UUID, req *CreateTaskCallbackRequest) (*TaskCallbackSubscriptionResponse, error)
- func (s *Service) DeleteTaskCallbackSubscription(ctx context.Context, runID, subscriptionID, userID uuid.UUID) error
- func (s *Service) EnqueueDelivery(ctx context.Context, run *db.Run, agentSlug string, ...) error
- func (s *Service) EnqueueRunEvent(ctx context.Context, event db.RunEvent) error
- func (s *Service) EnqueueRunEventDurable(ctx context.Context, event db.RunEvent) error
- func (s *Service) ListDeliveries(ctx context.Context, agentID, userID uuid.UUID, limit int) ([]DeliveryListItem, error)
- func (s *Service) ListTaskCallbackDeliveries(ctx context.Context, runID, userID uuid.UUID, limit int) ([]TaskCallbackDeliveryResponse, error)
- func (s *Service) ListTaskCallbackSubscriptions(ctx context.Context, runID, userID uuid.UUID) ([]TaskCallbackSubscriptionResponse, error)
- func (s *Service) ListTaskCallbackSubscriptionsForOwner(ctx context.Context, userID uuid.UUID, status string, limit int) ([]TaskCallbackSubscriptionResponse, error)
- func (s *Service) ResetWebhookEffectDelivery(ctx context.Context, effect db.RunEffectOutbox) error
- func (s *Service) RotateSecret(ctx context.Context, agentID, userID uuid.UUID) (*SetWebhookResponse, error)
- func (s *Service) SetWebhook(ctx context.Context, agentID, userID uuid.UUID, url string) (*SetWebhookResponse, error)
- func (s *Service) UpdateTaskCallbackSubscriptionStatus(ctx context.Context, runID, subscriptionID, userID uuid.UUID, status string) (*TaskCallbackSubscriptionResponse, error)
- type SetWebhookRequest
- type SetWebhookResponse
- type TaskCallbackDeliveryListResponse
- type TaskCallbackDeliveryResponse
- type TaskCallbackPayload
- type TaskCallbackSubscriptionListResponse
- type TaskCallbackSubscriptionResponse
- type WebhookPayload
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func StartWorker ¶
StartWorker 启动调用方 task callback 后台重试 worker(main.go 调,应该在 goroutine 中)。
行为:
- 每 workerInterval 扫一次 task_callback_deliveries:status='pending' AND next_retry_at <= NOW()
- 逐条调 AttemptTaskCallbackDelivery(顺序处理,避免对单一调用方 endpoint 并发打压)
- ctx.Done() → 优雅退出
Types ¶
type BatchTaskCallbackSubscriptionsResponse ¶
type BatchTaskCallbackSubscriptionsResponse struct {
Action string `json:"action"`
UpdatedCount int `json:"updated_count"`
Items []TaskCallbackSubscriptionResponse `json:"items"`
}
type CreateTaskCallbackRequest ¶
type CreateTaskCallbackRequest struct {
URL string `json:"target_url" validate:"required,url,max=500"`
Secret string `json:"secret,omitempty" validate:"omitempty,max=1000"`
EventTypes []string `` /* 287-byte string literal not displayed */
AuthScheme string `json:"auth_scheme,omitempty" validate:"omitempty,max=80"`
AuthCredentials string `json:"auth_credentials,omitempty" validate:"omitempty,max=1000"`
Metadata map[string]interface{} `json:"metadata,omitempty"`
}
CreateTaskCallbackRequest POST /api/v1/runs/:id/task-callbacks.
type DeliveryListItem ¶
type DeliveryListItem struct {
ID string `json:"id"`
RunID string `json:"run_id"`
URL string `json:"url"`
Status string `json:"status"`
ResponseStatus *int32 `json:"response_status,omitempty"`
ErrorMessage *string `json:"error_message,omitempty"`
AttemptCount int32 `json:"attempt_count"`
NextRetryAt *string `json:"next_retry_at,omitempty"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
DeliveryListItem describes removed Agent webhook delivery rows.
payload / response_body 不返回(避免巨大响应;如需 debug 可再加详情接口)。
type Handler ¶
type Handler struct {
// contains filtered or unexported fields
}
Handler webhook HTTP 入口(创作者侧)。
func NewHandler ¶
NewHandler 构造 Handler。cfg 可选(保持与其它模块一致)。
func (*Handler) BatchManageTaskCallbacks ¶
func (*Handler) CreateTaskCallback ¶
CreateTaskCallback 为单个 run 注册调用方任务回调。secret 仅本次返回。
func (*Handler) ListManagedTaskCallbacks ¶
func (*Handler) ListTaskCallbackDeliveries ¶
func (*Handler) RegisterProtected ¶
func (h *Handler) RegisterProtected(api *echo.Group, jwtMiddleware echo.MiddlewareFunc)
RegisterProtected 注册创作者侧端点(需 JWT)。
POST /runs/:id/task-callbacks 为单个 run 注册任务回调 GET /runs/:id/task-callbacks 查看 run 任务回调 GET /runs/:id/task-callbacks/deliveries 查看 run 任务回调投递记录 POST /runs/:id/task-callbacks/:callbackID/pause 暂停 run 任务回调 POST /runs/:id/task-callbacks/:callbackID/resume 恢复 run 任务回调 DELETE /runs/:id/task-callbacks/:callbackID 删除 run 任务回调 GET /task-callbacks 汇总当前用户的任务回调 POST /task-callbacks/batch 批量 pause / resume / delete
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
Service webhook 业务逻辑层。
关键约束(见 docs/13 子轮 2.1):
- 投递不阻塞 /run 响应(runtime 用 goroutine 调 EnqueueDelivery)
- HTTP 超时 15s
- 失败按 1min / 5min / 30min 三次重试,第 3 次失败终态 failed
- 签名 HMAC-SHA256,header X-OpenLinker-Signature: sha256=<hex>
func NewService ¶
NewService 构造 Service。
func (*Service) AttemptAgentWebhookEffect ¶ added in v0.1.56
func (s *Service) AttemptAgentWebhookEffect( ctx context.Context, effect db.RunEffectOutbox, ) runtimepkg.RunEffectAttemptResult
func (*Service) AttemptDelivery ¶
AttemptDelivery 单次投递。返回 nil 表示已正确处理 DB 状态(成功 / 已写重试时间 / 已 final)。
流程:
- 取 delivery(含 secret)
- 若已 success/failed 则跳过(worker 并发兜底)
- 若 secret 为 NULL(agent 已 ClearWebhook),直接 final failed
- HTTP POST 带签名
- 根据 HTTP 结果决定 success / retry / final fail
func (*Service) AttemptTaskCallbackDelivery ¶
AttemptTaskCallbackDelivery performs one task callback delivery attempt.
func (*Service) AttemptTaskCallbackEffect ¶ added in v0.1.56
func (s *Service) AttemptTaskCallbackEffect( ctx context.Context, effect db.RunEffectOutbox, ) runtimepkg.RunEffectAttemptResult
func (*Service) BatchManageTaskCallbackSubscriptions ¶
func (s *Service) BatchManageTaskCallbackSubscriptions(ctx context.Context, userID uuid.UUID, req *BatchTaskCallbackSubscriptionsRequest) (*BatchTaskCallbackSubscriptionsResponse, error)
func (*Service) ClearWebhook ¶
ClearWebhook 清除 webhook 配置(url 与 secret 都置 NULL)。
func (*Service) CreateTaskCallbackSubscription ¶
func (s *Service) CreateTaskCallbackSubscription(ctx context.Context, runID, userID uuid.UUID, req *CreateTaskCallbackRequest) (*TaskCallbackSubscriptionResponse, error)
CreateTaskCallbackSubscription registers a signed caller-owned callback for one run.
func (*Service) DeleteTaskCallbackSubscription ¶
func (s *Service) DeleteTaskCallbackSubscription(ctx context.Context, runID, subscriptionID, userID uuid.UUID) error
DeleteTaskCallbackSubscription soft-deletes a task callback.
func (*Service) EnqueueDelivery ¶
func (s *Service) EnqueueDelivery(ctx context.Context, run *db.Run, agentSlug string, output map[string]interface{}) error
EnqueueDelivery handles the legacy Agent webhook queue.
流程:
- 查 legacy Agent webhook URL:NULL → 直接 return(不投递)
- 构造 payload(event=run.completed)
- INSERT webhook_deliveries (status=pending, next_retry_at=NOW())
- 立即异步触发第一次投递(不等待)
func (*Service) EnqueueRunEvent ¶
EnqueueRunEvent creates deliveries for all subscriptions interested in this run_event.
func (*Service) EnqueueRunEventDurable ¶ added in v0.1.56
EnqueueRunEventDurable bypasses the negative subscription cache and returns every materialization error to its caller. The parent-completion Effect uses this path so it cannot be marked succeeded until all matching callback delivery rows exist durably.
func (*Service) ListDeliveries ¶
func (s *Service) ListDeliveries(ctx context.Context, agentID, userID uuid.UUID, limit int) ([]DeliveryListItem, error)
ListDeliveries 创作者查看 agent 投递历史。
func (*Service) ListTaskCallbackDeliveries ¶
func (*Service) ListTaskCallbackSubscriptions ¶
func (s *Service) ListTaskCallbackSubscriptions(ctx context.Context, runID, userID uuid.UUID) ([]TaskCallbackSubscriptionResponse, error)
ListTaskCallbackSubscriptions returns active/non-deleted task callbacks.
func (*Service) ListTaskCallbackSubscriptionsForOwner ¶
func (*Service) ResetWebhookEffectDelivery ¶ added in v0.1.56
func (*Service) RotateSecret ¶
func (s *Service) RotateSecret(ctx context.Context, agentID, userID uuid.UUID) (*SetWebhookResponse, error)
RotateSecret 重新生成 secret,保留原 url。
agent 必须已配 webhook_url;否则 404。
func (*Service) SetWebhook ¶
func (s *Service) SetWebhook(ctx context.Context, agentID, userID uuid.UUID, url string) (*SetWebhookResponse, error)
SetWebhook 创作者设置 webhook_url,生成新 secret 并返回(仅本次返回)。
校验:
- agent 必须存在
- agent.creator_id == userID(防越权)
- URL 必须满足共享出网策略(schema CHECK 兜底,但前置返回友好错误)
func (*Service) UpdateTaskCallbackSubscriptionStatus ¶
func (s *Service) UpdateTaskCallbackSubscriptionStatus(ctx context.Context, runID, subscriptionID, userID uuid.UUID, status string) (*TaskCallbackSubscriptionResponse, error)
UpdateTaskCallbackSubscriptionStatus pauses or resumes a task callback.
type SetWebhookRequest ¶
type SetWebhookRequest struct {
URL string `json:"webhook_url" validate:"required,url,startswith=https://,max=500"`
}
SetWebhookRequest is kept for removed Agent webhook route tests.
URL 必须 https,由 schema CHECK 兜底(agents_webhook_https),同时这里前置校验。
type SetWebhookResponse ¶
type SetWebhookResponse struct {
URL string `json:"webhook_url"`
Secret string `json:"webhook_secret"`
}
SetWebhookResponse 创建 / 重置 webhook 响应。
Secret 仅在创建 / rotate 时返回一次,后续无法再取(DB 仅存原值,但接口不再返回)。
type TaskCallbackDeliveryListResponse ¶
type TaskCallbackDeliveryListResponse struct {
Items []TaskCallbackDeliveryResponse `json:"items"`
}
type TaskCallbackDeliveryResponse ¶
type TaskCallbackDeliveryResponse struct {
ID string `json:"id"`
SubscriptionID string `json:"subscription_id"`
RunEventID string `json:"run_event_id"`
EventType string `json:"event_type"`
TargetURL string `json:"target_url"`
Status string `json:"status"`
ResponseStatus *int32 `json:"response_status,omitempty"`
ErrorMessage *string `json:"error_message,omitempty"`
AttemptCount int32 `json:"attempt_count"`
NextRetryAt *string `json:"next_retry_at,omitempty"`
DeliveredAt *string `json:"delivered_at,omitempty"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
TaskCallbackDeliveryResponse describes one signed callback delivery attempt.
type TaskCallbackPayload ¶
type TaskCallbackPayload struct {
EventID string `json:"event_id"`
RunID string `json:"run_id"`
ParentRunID string `json:"parent_run_id,omitempty"`
EventType string `json:"event_type"`
Sequence int32 `json:"sequence"`
Payload map[string]interface{} `json:"payload"`
SubscriptionID string `json:"subscription_id"`
CreatedAt string `json:"created_at"`
// contains filtered or unexported fields
}
TaskCallbackPayload is the A2A-style push body derived from run_events.
func (TaskCallbackPayload) MarshalJSON ¶
func (p TaskCallbackPayload) MarshalJSON() ([]byte, error)
type TaskCallbackSubscriptionListResponse ¶
type TaskCallbackSubscriptionListResponse struct {
Items []TaskCallbackSubscriptionResponse `json:"items"`
}
type TaskCallbackSubscriptionResponse ¶
type TaskCallbackSubscriptionResponse struct {
ID string `json:"id"`
RunID string `json:"run_id"`
TargetURL string `json:"target_url"`
EventTypes []string `json:"event_types"`
AuthScheme string `json:"auth_scheme,omitempty"`
Status string `json:"status"`
ConsecutiveFailures int32 `json:"consecutive_failures"`
Secret string `json:"secret,omitempty"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
TaskCallbackSubscriptionResponse is returned to the run owner.
type WebhookPayload ¶
type WebhookPayload struct {
Event string `json:"event"`
RunID string `json:"run_id"`
AgentID string `json:"agent_id"`
AgentSlug string `json:"agent_slug"`
UserID string `json:"user_id"`
Status string `json:"status"`
Input map[string]interface{} `json:"input"`
Output map[string]interface{} `json:"output,omitempty"`
ErrorCode string `json:"error_code,omitempty"`
ErrorMessage string `json:"error_message,omitempty"`
CostCents int32 `json:"cost_cents"`
DurationMs int32 `json:"duration_ms"`
StartedAt string `json:"started_at"`
FinishedAt string `json:"finished_at"`
}
WebhookPayload is the legacy Agent webhook body.
字段含义:
- Event 固定 "run.completed"(成功 / 失败 / 超时 都是该事件,由 status 区分)
- Status: 'success' / 'failed' / 'timeout'
- Output 仅 success 时非空
- ErrorCode / ErrorMessage 仅 failed / timeout 时非空
- CostCents:成功时 = agent 价格;失败 / 超时 = 0(已退款)