workflow

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

View Source
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 CancellationEvidence struct {
	CancellationID uuid.UUID
	WorkflowRunID  uuid.UUID
	ActorUserID    uuid.UUID
	ReasonCode     string
	State          string
	ErrorCode      string
	RequestedAt    time.Time
	AppliedAt      *time.Time
	FinishedAt     *time.Time
}

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 ExternalExecutionLaunchFence struct {
	CallerServiceID   string
	ExternalRequestID uuid.UUID
	ActorUserID       uuid.UUID
	LaunchToken       uuid.UUID
}

type ExternalExecutionTargetValidation added in v0.1.56

type ExternalExecutionTargetValidation struct {
	TargetName        string `json:"target_name"`
	Executable        bool   `json:"executable"`
	UnavailableReason string `json:"unavailable_reason,omitempty"`
	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) CancelRun

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

func (*Handler) CompareRuns

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

func (*Handler) Create

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

func (*Handler) Get

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

func (*Handler) GetRun

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

func (*Handler) List

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

func (*Handler) ListRuns

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

func (*Handler) PauseRun

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

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

func (*Handler) RerunStep

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

func (*Handler) ResumeRun

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

func (*Handler) RetryRun

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

func (*Handler) Run

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

func (*Handler) StartRun

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

type RerunWorkflowStepRequest

type RerunWorkflowStepRequest struct {
	NodeKey string `json:"node_key" validate:"required,min=1,max=80"`
}

type RunWorkerConfig

type RunWorkerConfig struct {
	Interval   time.Duration
	StaleAfter time.Duration
	ClaimBurst int
}

type RunWorkflowRequest

type RunWorkflowRequest struct {
	Input       map[string]interface{} `json:"input,omitempty"`
	MaxAttempts int32                  `json:"max_attempts,omitempty" validate:"omitempty,min=1,max=10"`
}

type Service

type Service struct {
	// contains filtered or unexported fields
}

Service persists workflows and executes Agent nodes through the runtime service.

func NewService

func NewService(pool *pgxpool.Pool, runtimeSvc *runtime.Service) *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 (s *Service) CancelWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)

func (*Service) ClaimAndRunPendingWorkflow

func (s *Service) ClaimAndRunPendingWorkflow(ctx context.Context) (bool, error)

func (*Service) CompareWorkflowRuns

func (s *Service) CompareWorkflowRuns(ctx context.Context, userID, baseRunID, candidateRunID uuid.UUID) (*WorkflowRunComparisonResponse, error)

func (*Service) CreateWorkflow

func (s *Service) CreateWorkflow(ctx context.Context, userID uuid.UUID, req *CreateWorkflowRequest) (*WorkflowResponse, error)

func (*Service) GetWorkflow

func (s *Service) GetWorkflow(ctx context.Context, userID, workflowID uuid.UUID) (*WorkflowResponse, error)

func (*Service) GetWorkflowCancellationEvidence added in v0.1.56

func (s *Service) GetWorkflowCancellationEvidence(
	ctx context.Context,
	actorUserID, workflowRunID uuid.UUID,
) (CancellationEvidence, error)

func (*Service) GetWorkflowRun

func (s *Service) GetWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)

func (*Service) ListWorkflowRuns

func (s *Service) ListWorkflowRuns(ctx context.Context, userID, workflowID uuid.UUID, limit int32) (*WorkflowRunListResponse, error)

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 (s *Service) ListWorkflows(ctx context.Context, userID uuid.UUID, limit int32) (*WorkflowListResponse, error)

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 (s *Service) PauseWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)

func (*Service) RequeueStaleWorkflowRuns

func (s *Service) RequeueStaleWorkflowRuns(ctx context.Context, staleAfter time.Duration) (int64, error)

func (*Service) RerunWorkflowStep

func (s *Service) RerunWorkflowStep(ctx context.Context, userID, workflowRunID uuid.UUID, req *RerunWorkflowStepRequest) (*WorkflowStepRerunResponse, error)

func (*Service) ResumeWorkflowRun

func (s *Service) ResumeWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)

func (*Service) RetryWorkflowRun

func (s *Service) RetryWorkflowRun(ctx context.Context, userID, workflowRunID uuid.UUID) (*WorkflowRunResponse, error)

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 WorkflowListResponse struct {
	Items        []WorkflowResponse `json:"items"`
	Total        int32              `json:"total"`
	Page         int32              `json:"page"`
	Size         int32              `json:"size"`
	Query        string             `json:"query,omitempty"`
	Sort         string             `json:"sort"`
	StatusFilter string             `json:"status_filter,omitempty"`
}

type WorkflowNodeRequest

type WorkflowNodeRequest struct {
	Key     string                 `json:"key" validate:"required,min=1,max=80"`
	Title   string                 `json:"title,omitempty" validate:"omitempty,max=160"`
	AgentID uuid.UUID              `json:"agent_id" validate:"required"`
	Config  map[string]interface{} `json:"config,omitempty"`
}

type WorkflowNodeResponse

type WorkflowNodeResponse struct {
	ID       string                 `json:"id"`
	Key      string                 `json:"key"`
	Type     string                 `json:"type"`
	AgentID  string                 `json:"agent_id"`
	Title    string                 `json:"title"`
	Config   map[string]interface{} `json:"config"`
	Position int32                  `json:"position"`
}

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 WorkflowRunListResponse struct {
	Items        []WorkflowRunResponse `json:"items"`
	Total        int32                 `json:"total"`
	Page         int32                 `json:"page"`
	Size         int32                 `json:"size"`
	Query        string                `json:"query,omitempty"`
	Sort         string                `json:"sort"`
	StatusFilter string                `json:"status_filter,omitempty"`
}

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"`
}

Jump to

Keyboard shortcuts

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