Documentation
¶
Index ¶
- Variables
- func StartRunWorker(ctx context.Context, svc *Service, cfg RunWorkerConfig)
- type CancellationEvidence
- type CreateWorkflowRequest
- type ExternalExecutionLaunchFence
- type ExternalExecutionTargetValidation
- type Handler
- func (h *Handler) CancelRun(c echo.Context) error
- func (h *Handler) CompareRuns(c echo.Context) error
- func (h *Handler) Create(c echo.Context) error
- func (h *Handler) Get(c echo.Context) error
- func (h *Handler) GetRun(c echo.Context) error
- func (h *Handler) List(c echo.Context) error
- func (h *Handler) ListRuns(c echo.Context) error
- func (h *Handler) PauseRun(c echo.Context) error
- func (h *Handler) RegisterProtected(api *echo.Group, jwtMiddleware echo.MiddlewareFunc)
- func (h *Handler) RerunStep(c echo.Context) error
- func (h *Handler) ResumeRun(c echo.Context) error
- func (h *Handler) RetryRun(c echo.Context) error
- func (h *Handler) Run(c echo.Context) error
- func (h *Handler) StartRun(c echo.Context) error
- type RerunWorkflowStepRequest
- type RunWorkerConfig
- type RunWorkflowRequest
- type Service
- func (s *Service) CancelExternalWorkflowRun(ctx context.Context, actorUserID, workflowRunID uuid.UUID, reasonCode string) (*WorkflowRunResponse, CancellationEvidence, error)
- func (s *Service) CancelWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)
- func (s *Service) ClaimAndRunPendingWorkflow(ctx context.Context) (bool, error)
- func (s *Service) CompareWorkflowRuns(ctx context.Context, userID, baseRunID, candidateRunID uuid.UUID) (*WorkflowRunComparisonResponse, error)
- func (s *Service) CreateWorkflow(ctx context.Context, userID uuid.UUID, req *CreateWorkflowRequest) (*WorkflowResponse, error)
- func (s *Service) GetWorkflow(ctx context.Context, userID, workflowID uuid.UUID) (*WorkflowResponse, error)
- func (s *Service) GetWorkflowCancellationEvidence(ctx context.Context, actorUserID, workflowRunID uuid.UUID) (CancellationEvidence, error)
- func (s *Service) GetWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)
- func (s *Service) ListWorkflowRuns(ctx context.Context, userID, workflowID uuid.UUID, limit int32) (*WorkflowRunListResponse, error)
- func (s *Service) ListWorkflowRunsPage(ctx context.Context, userID, workflowID uuid.UUID, query, status, sort string, ...) (*WorkflowRunListResponse, error)
- func (s *Service) ListWorkflows(ctx context.Context, userID uuid.UUID, limit int32) (*WorkflowListResponse, error)
- func (s *Service) ListWorkflowsPage(ctx context.Context, userID uuid.UUID, query, status, sort string, ...) (*WorkflowListResponse, error)
- func (s *Service) LookupExternalExecutionWorkflowRun(ctx context.Context, callerServiceID string, ...) (*WorkflowRunResponse, bool, error)
- func (s *Service) LookupExternalExecutionWorkflowRunByIdentity(ctx context.Context, callerServiceID string, ...) (*WorkflowRunResponse, bool, error)
- func (s *Service) PauseWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)
- func (s *Service) RequeueStaleWorkflowRuns(ctx context.Context, staleAfter time.Duration) (int64, error)
- func (s *Service) RerunWorkflowStep(ctx context.Context, userID, workflowRunID uuid.UUID, ...) (*WorkflowStepRerunResponse, error)
- func (s *Service) ResumeWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)
- func (s *Service) RetryWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)
- func (s *Service) RunWorkflow(ctx context.Context, userID, workflowID uuid.UUID, req *RunWorkflowRequest) (*WorkflowRunResponse, error)
- func (s *Service) SetRunUpdateSource(source runtime.RunUpdateSource)
- func (s *Service) SetWorkerObserver(observer runtime.WorkerObserver)
- func (s *Service) StartExternalExecutionWorkflowRun(ctx context.Context, callerServiceID string, ...) (*WorkflowRunResponse, error)
- func (s *Service) StartExternalExecutionWorkflowRunWithFence(ctx context.Context, targetOwnerID, actorUserID, workflowID uuid.UUID, ...) (*WorkflowRunResponse, error)
- func (s *Service) StartWorkflowRun(ctx context.Context, userID, workflowID uuid.UUID, req *RunWorkflowRequest) (*WorkflowRunResponse, error)
- func (s *Service) ValidateExternalExecutionTarget(ctx context.Context, targetOwnerID, workflowID uuid.UUID) (*ExternalExecutionTargetValidation, error)
- type WorkflowListResponse
- type WorkflowNodeRequest
- type WorkflowNodeResponse
- type WorkflowResponse
- type WorkflowRunComparisonResponse
- type WorkflowRunListResponse
- type WorkflowRunResponse
- type WorkflowRunStepCompareResponse
- type WorkflowRunStepResponse
- type WorkflowStepRerunResponse
Constants ¶
This section is empty.
Variables ¶
var ErrExternalWorkflowLaunchFenceRejected = errors.New("external workflow launch fence rejected")
Functions ¶
func StartRunWorker ¶
func StartRunWorker(ctx context.Context, svc *Service, cfg RunWorkerConfig)
Types ¶
type CancellationEvidence ¶ added in v0.1.56
type CreateWorkflowRequest ¶
type CreateWorkflowRequest struct {
Name string `json:"name" validate:"required,min=1,max=120"`
Description string `json:"description,omitempty" validate:"omitempty,max=500"`
Nodes []WorkflowNodeRequest `json:"nodes" validate:"required,min=1,max=10,dive"`
Edges []map[string]interface{} `json:"edges" validate:"omitempty,max=20"`
}
CreateWorkflowRequest creates a DAG of Agent nodes. Omitted/null edges use the sequential default; an explicit empty array makes all nodes independent.
type ExternalExecutionLaunchFence ¶ added in v0.1.56
type ExternalExecutionTargetValidation ¶ added in v0.1.56
type ExternalExecutionTargetValidation struct {
TargetName string `json:"target_name"`
Executable bool `json:"executable"`
ContractHash string `json:"contract_hash,omitempty"`
}
ExternalExecutionTargetValidation is the narrow Core result consumed by an external execution caller. It deliberately carries no workflow definition or endpoint details.
type Handler ¶
type Handler struct {
// contains filtered or unexported fields
}
func NewHandler ¶
func NewHandler(svc workflowService) *Handler
func (*Handler) RegisterProtected ¶
func (h *Handler) RegisterProtected(api *echo.Group, jwtMiddleware echo.MiddlewareFunc)
RegisterProtected mounts workflow APIs behind user auth.
POST /workflows 创建 workflow GET /workflows 查询自己的 workflow 列表 GET /workflows/:id 查询 workflow 定义 POST /workflows/:id/run 同步执行 workflow POST /workflows/:id/runs 异步启动 workflow run GET /workflows/:id/runs 查询 workflow run 历史 GET /workflow-runs/:id 查询 workflow run POST /workflow-runs/:id/retry 复制失败 run 输入并重新入队 POST /workflow-runs/:id/steps/rerun 基于既有 run 重跑某个 step 及其下游 GET /workflow-runs/:id/compare/:other_id 对比两个 workflow run POST /workflow-runs/:id/pause 暂停 pending/running run POST /workflow-runs/:id/resume 恢复 paused run POST /workflow-runs/:id/cancel 取消 pending/running/paused run
type RerunWorkflowStepRequest ¶
type RerunWorkflowStepRequest struct {
NodeKey string `json:"node_key" validate:"required,min=1,max=80"`
}
type RunWorkerConfig ¶
type RunWorkflowRequest ¶
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
Service persists workflows and executes Agent nodes through the runtime service.
func (*Service) CancelExternalWorkflowRun ¶ added in v0.1.56
func (s *Service) CancelExternalWorkflowRun( ctx context.Context, actorUserID, workflowRunID uuid.UUID, reasonCode string, ) (*WorkflowRunResponse, CancellationEvidence, error)
CancelExternalWorkflowRun persists first-writer intent, prevents all new child launches, then reconciles physical child stop evidence.
func (*Service) CancelWorkflowRun ¶
func (*Service) ClaimAndRunPendingWorkflow ¶
func (*Service) CompareWorkflowRuns ¶
func (*Service) CreateWorkflow ¶
func (s *Service) CreateWorkflow(ctx context.Context, userID uuid.UUID, req *CreateWorkflowRequest) (*WorkflowResponse, error)
func (*Service) GetWorkflow ¶
func (*Service) GetWorkflowCancellationEvidence ¶ added in v0.1.56
func (*Service) GetWorkflowRun ¶
func (*Service) ListWorkflowRuns ¶
func (*Service) ListWorkflowRunsPage ¶ added in v0.1.18
func (s *Service) ListWorkflowRunsPage(ctx context.Context, userID, workflowID uuid.UUID, query, status, sort string, page, size int32) (*WorkflowRunListResponse, error)
ListWorkflowRunsPage returns workflow run history with server-side search, status filtering, sorting, and pagination.
func (*Service) ListWorkflows ¶
func (*Service) ListWorkflowsPage ¶ added in v0.1.15
func (s *Service) ListWorkflowsPage(ctx context.Context, userID uuid.UUID, query, status, sort string, page, size int32) (*WorkflowListResponse, error)
ListWorkflowsPage returns workflows with server-side search, status filter, sort, and pagination.
func (*Service) LookupExternalExecutionWorkflowRun ¶ added in v0.1.56
func (s *Service) LookupExternalExecutionWorkflowRun( ctx context.Context, callerServiceID string, actorUserID, workflowID, externalRequestID uuid.UUID, input map[string]interface{}, ) (*WorkflowRunResponse, bool, error)
LookupExternalExecutionWorkflowRun performs the read-only half of external workflow execution idempotency. The deterministic caller/request identity is resolved without consulting the mutable Workflow definition. A committed row must still match every execution semantic; a missing row is returned as a clean miss and never creates a workflow run.
func (*Service) LookupExternalExecutionWorkflowRunByIdentity ¶ added in v0.1.56
func (s *Service) LookupExternalExecutionWorkflowRunByIdentity( ctx context.Context, callerServiceID string, actorUserID, workflowID, externalRequestID uuid.UUID, ) (*WorkflowRunResponse, bool, error)
LookupExternalExecutionWorkflowRunByIdentity is the key-only recovery path. The deterministic ID plus actor/target checks are sufficient; no Cloud input is required and the method never mutates state.
func (*Service) PauseWorkflowRun ¶
func (*Service) RequeueStaleWorkflowRuns ¶
func (*Service) RerunWorkflowStep ¶
func (s *Service) RerunWorkflowStep(ctx context.Context, userID, workflowRunID uuid.UUID, req *RerunWorkflowStepRequest) (*WorkflowStepRerunResponse, error)
func (*Service) ResumeWorkflowRun ¶
func (*Service) RetryWorkflowRun ¶
func (*Service) RunWorkflow ¶
func (s *Service) RunWorkflow(ctx context.Context, userID, workflowID uuid.UUID, req *RunWorkflowRequest) (*WorkflowRunResponse, error)
func (*Service) SetRunUpdateSource ¶ added in v0.1.56
func (s *Service) SetRunUpdateSource(source runtime.RunUpdateSource)
SetRunUpdateSource enables advisory event-driven child Run waits. The Run row remains authoritative and the legacy polling path remains the degraded fallback when the shared listener is unavailable.
func (*Service) SetWorkerObserver ¶ added in v0.1.56
func (s *Service) SetWorkerObserver(observer runtime.WorkerObserver)
SetWorkerObserver installs payload-free test instrumentation only.
func (*Service) StartExternalExecutionWorkflowRun ¶ added in v0.1.56
func (s *Service) StartExternalExecutionWorkflowRun( ctx context.Context, callerServiceID string, targetOwnerID, actorUserID, workflowID, externalRequestID uuid.UUID, input map[string]interface{}, ) (*WorkflowRunResponse, error)
StartExternalExecutionWorkflowRun keeps target ownership and result ownership separate. The run ID is derived from the verified caller plus its request ID, so retries remain idempotent without allowing cross-service collisions.
func (*Service) StartExternalExecutionWorkflowRunWithFence ¶ added in v0.1.56
func (s *Service) StartExternalExecutionWorkflowRunWithFence( ctx context.Context, targetOwnerID, actorUserID, workflowID uuid.UUID, input map[string]interface{}, fence ExternalExecutionLaunchFence, ) (*WorkflowRunResponse, error)
StartExternalExecutionWorkflowRunWithFence commits the deterministic Workflow Run only while the external launch token is current. The external key row is always locked before the external execution row.
func (*Service) StartWorkflowRun ¶
func (s *Service) StartWorkflowRun(ctx context.Context, userID, workflowID uuid.UUID, req *RunWorkflowRequest) (*WorkflowRunResponse, error)
StartWorkflowRun creates a durable pending workflow run. A background worker will claim and execute it, so HTTP clients do not have to keep long requests open.
func (*Service) ValidateExternalExecutionTarget ¶ added in v0.1.56
func (s *Service) ValidateExternalExecutionTarget(ctx context.Context, targetOwnerID, workflowID uuid.UUID) (*ExternalExecutionTargetValidation, error)
ValidateExternalExecutionTarget verifies that an actor-owned workflow is safe for external execution. Passing uuid.Nil to the existing Agent validator intentionally requires every node to be public and callable.
type WorkflowListResponse ¶
type WorkflowNodeRequest ¶
type WorkflowNodeResponse ¶
type WorkflowResponse ¶
type WorkflowResponse struct {
ID string `json:"id"`
Name string `json:"name"`
Description string `json:"description"`
Status string `json:"status"`
Nodes []WorkflowNodeResponse `json:"nodes"`
Edges []map[string]interface{} `json:"edges"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
type WorkflowRunComparisonResponse ¶
type WorkflowRunComparisonResponse struct {
BaseRunID string `json:"base_run_id"`
CandidateRunID string `json:"candidate_run_id"`
WorkflowID string `json:"workflow_id"`
StatusChanged bool `json:"status_changed"`
OutputChanged bool `json:"output_changed"`
ChangedNodeKeys []string `json:"changed_node_keys"`
Steps []WorkflowRunStepCompareResponse `json:"steps"`
}
type WorkflowRunListResponse ¶
type WorkflowRunResponse ¶
type WorkflowRunResponse struct {
ID string `json:"id"`
WorkflowID string `json:"workflow_id"`
Status string `json:"status"`
Input map[string]interface{} `json:"input"`
Output map[string]interface{} `json:"output,omitempty"`
Error string `json:"error_message,omitempty"`
Steps []WorkflowRunStepResponse `json:"steps"`
AttemptCount int32 `json:"attempt_count"`
MaxAttempts int32 `json:"max_attempts"`
NextRetryAt string `json:"next_retry_at,omitempty"`
ClaimedAt string `json:"claimed_at,omitempty"`
LastWorkerError string `json:"last_worker_error,omitempty"`
StartedAt string `json:"started_at"`
FinishedAt string `json:"finished_at,omitempty"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
type WorkflowRunStepCompareResponse ¶
type WorkflowRunStepCompareResponse struct {
NodeKey string `json:"node_key"`
BaseStatus string `json:"base_status,omitempty"`
CandidateStatus string `json:"candidate_status,omitempty"`
BaseRunID string `json:"base_run_id,omitempty"`
CandidateRunID string `json:"candidate_run_id,omitempty"`
StatusChanged bool `json:"status_changed"`
RunChanged bool `json:"run_changed"`
OutputChanged bool `json:"output_changed"`
ErrorChanged bool `json:"error_changed"`
Changed bool `json:"changed"`
}
type WorkflowRunStepResponse ¶
type WorkflowRunStepResponse struct {
ID string `json:"id"`
NodeID string `json:"node_id"`
NodeKey string `json:"node_key"`
AgentID string `json:"agent_id"`
RunID string `json:"run_id,omitempty"`
Status string `json:"status"`
Input map[string]interface{} `json:"input"`
Output map[string]interface{} `json:"output,omitempty"`
Error string `json:"error_message,omitempty"`
Sequence int32 `json:"sequence"`
StartedAt string `json:"started_at"`
FinishedAt string `json:"finished_at,omitempty"`
}
type WorkflowStepRerunResponse ¶
type WorkflowStepRerunResponse struct {
SourceRunID string `json:"source_run_id"`
RerunRunID string `json:"rerun_run_id"`
NodeKey string `json:"node_key"`
ReusedNodeKeys []string `json:"reused_node_keys"`
RerunNodeKeys []string `json:"rerun_node_keys"`
Run WorkflowRunResponse `json:"run"`
Comparison WorkflowRunComparisonResponse `json:"comparison"`
}