webhook

package
v0.1.59 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 14, 2026 License: Apache-2.0 Imports: 27 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func StartWorker

func StartWorker(ctx context.Context, svc *Service)

StartWorker 启动调用方 task callback 后台重试 worker(main.go 调,应该在 goroutine 中)。

行为:

  • 每 workerInterval 扫一次 task_callback_deliveries:status='pending' AND next_retry_at <= NOW()
  • 逐条调 AttemptTaskCallbackDelivery(顺序处理,避免对单一调用方 endpoint 并发打压)
  • ctx.Done() → 优雅退出

Types

type BatchTaskCallbackSubscriptionsRequest

type BatchTaskCallbackSubscriptionsRequest struct {
	SubscriptionIDs []string `json:"subscription_ids" validate:"required,min=1,max=50,dive,uuid"`
	Action          string   `json:"action" validate:"required,oneof=pause resume delete"`
}

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

func NewHandler(svc webhookService, cfg ...*config.Config) *Handler

NewHandler 构造 Handler。cfg 可选(保持与其它模块一致)。

func (*Handler) BatchManageTaskCallbacks

func (h *Handler) BatchManageTaskCallbacks(c echo.Context) error

func (*Handler) CreateTaskCallback

func (h *Handler) CreateTaskCallback(c echo.Context) error

CreateTaskCallback 为单个 run 注册调用方任务回调。secret 仅本次返回。

func (*Handler) DeleteTaskCallback

func (h *Handler) DeleteTaskCallback(c echo.Context) error

func (*Handler) ListManagedTaskCallbacks

func (h *Handler) ListManagedTaskCallbacks(c echo.Context) error

func (*Handler) ListTaskCallbackDeliveries

func (h *Handler) ListTaskCallbackDeliveries(c echo.Context) error

func (*Handler) ListTaskCallbacks

func (h *Handler) ListTaskCallbacks(c echo.Context) error

func (*Handler) PauseTaskCallback

func (h *Handler) PauseTaskCallback(c echo.Context) error

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

func (*Handler) ResumeTaskCallback

func (h *Handler) ResumeTaskCallback(c echo.Context) error

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

func NewService(pool *pgxpool.Pool, cfg ...*config.Config) *Service

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

func (s *Service) AttemptDelivery(ctx context.Context, deliveryID uuid.UUID) error

AttemptDelivery 单次投递。返回 nil 表示已正确处理 DB 状态(成功 / 已写重试时间 / 已 final)。

流程:

  1. 取 delivery(含 secret)
  2. 若已 success/failed 则跳过(worker 并发兜底)
  3. 若 secret 为 NULL(agent 已 ClearWebhook),直接 final failed
  4. HTTP POST 带签名
  5. 根据 HTTP 结果决定 success / retry / final fail

func (*Service) AttemptTaskCallbackDelivery

func (s *Service) AttemptTaskCallbackDelivery(ctx context.Context, deliveryID uuid.UUID) error

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) ClearWebhook

func (s *Service) ClearWebhook(ctx context.Context, agentID, userID uuid.UUID) error

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.

流程:

  1. 查 legacy Agent webhook URL:NULL → 直接 return(不投递)
  2. 构造 payload(event=run.completed)
  3. INSERT webhook_deliveries (status=pending, next_retry_at=NOW())
  4. 立即异步触发第一次投递(不等待)

func (*Service) EnqueueRunEvent

func (s *Service) EnqueueRunEvent(ctx context.Context, event db.RunEvent) error

EnqueueRunEvent creates deliveries for all subscriptions interested in this run_event.

func (*Service) EnqueueRunEventDurable added in v0.1.56

func (s *Service) EnqueueRunEventDurable(ctx context.Context, event db.RunEvent) error

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 (s *Service) ListTaskCallbackDeliveries(ctx context.Context, runID, userID uuid.UUID, limit int) ([]TaskCallbackDeliveryResponse, error)

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 (s *Service) ListTaskCallbackSubscriptionsForOwner(ctx context.Context, userID uuid.UUID, status string, limit int) ([]TaskCallbackSubscriptionResponse, error)

func (*Service) ResetWebhookEffectDelivery added in v0.1.56

func (s *Service) ResetWebhookEffectDelivery(ctx context.Context, effect db.RunEffectOutbox) error

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 并返回(仅本次返回)。

校验:

  1. agent 必须存在
  2. agent.creator_id == userID(防越权)
  3. 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(已退款)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL