repository

package
v1.0.3 Latest Latest
Warning

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

Go to latest
Published: Jul 7, 2026 License: MIT Imports: 3 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type AwaitBindingRepository

type AwaitBindingRepository interface {
	Create(ctx context.Context, b *domain.AwaitBinding) error
	Update(ctx context.Context, b *domain.AwaitBinding) error

	GetByID(ctx context.Context, id int64) (*domain.AwaitBinding, error)
	ListByTaskID(ctx context.Context, taskID int64) ([]*domain.AwaitBinding, error)
	GetByTaskAndNode(ctx context.Context, taskID int64, nodeName string) (*domain.AwaitBinding, error)

	FindWaitingByProviderTaskID(ctx context.Context, provider, providerTaskID string) (*domain.AwaitBinding, error)
	FindWaitingByAPITaskID(ctx context.Context, provider, apiTaskID string) (*domain.AwaitBinding, error)
	FindWaitingBySignal(ctx context.Context, signalName, callbackToken string) (*domain.AwaitBinding, error)

	TransitionStatus(ctx context.Context, id int64, from domain.AwaitBindingStatus, to domain.AwaitBindingStatus) (bool, error)
	ClaimCompleting(ctx context.Context, id int64, expectedStatuses []domain.AwaitBindingStatus) (bool, error)

	FindPollDue(ctx context.Context, now time.Time, limit int) ([]*domain.AwaitBinding, error)
	FindTimeoutDue(ctx context.Context, now time.Time, limit int) ([]*domain.AwaitBinding, error)
}

type EventRepository

type EventRepository interface {
	Create(ctx context.Context, event *domain.TaskEvent) error
	FindByTaskID(ctx context.Context, taskID int64, isByRoot bool) ([]domain.TaskEvent, error)
	FindByTaskIDAndTypePrefixes(ctx context.Context, taskID int64, prefixes []string, isByRoot bool) ([]domain.TaskEvent, error)
	// FindPersistentByTaskID returns only persistent-grade events, ordered by sequence.
	// If afterSequence > 0, only events with sequence > afterSequence are returned (incremental recovery).
	FindPersistentByTaskID(ctx context.Context, taskID int64, afterSequence int64, limit int, isByRoot bool) ([]domain.TaskEvent, error)
}

type NodeRuntimeRepository

type NodeRuntimeRepository interface {
	Create(ctx context.Context, n *domain.NodeRuntime) error
	Update(ctx context.Context, n *domain.NodeRuntime) error
	// FindByTaskID 根据任务ID 查找所有节点
	FindByTaskID(ctx context.Context, taskID int64) ([]*domain.NodeRuntime, error)
	// FindByTaskIDAndNode 根据任务ID 和 节点名称 查找节点
	FindByTaskIDAndNode(ctx context.Context, taskID int64, node string) (*domain.NodeRuntime, error)
	// MarkRunningAsRetrying 标记正在运行中的节点为重试状态
	MarkRunningAsRetrying(ctx context.Context, taskID int64) error
	MarkAsRetrying(ctx context.Context, taskID int64, name string) error
	MarkFailed(ctx context.Context, taskID int64, name string, errMessage string) error

	FindExpiredRunningNodes(ctx context.Context, expire time.Time) ([]*domain.NodeRuntime, error)
	AttemptCompletePendingEdges(ctx context.Context, taskID int64, nodeName string, output map[string]any, errMsg string) (bool, error)

	CloneCheckpoint(ctx context.Context, fromTaskID, toTaskID int64) error
}

NodeRuntimeRepository 节点状态存储

type TaskCostTraceRepository

type TaskCostTraceRepository interface {
	Upsert(ctx context.Context, trace *domain.TaskCostTrace) error
	ListByTaskID(ctx context.Context, taskID int64) ([]*domain.TaskCostTrace, error)
	SumEstimatedCostByTaskID(ctx context.Context, taskID int64) (float64, error)
	SumEstimatedCostByRootTaskID(ctx context.Context, rootTaskID int64) (float64, error)
}

type TaskQueue

type TaskQueue interface {
	Push(ctx context.Context, taskID int64) error
	//Pop(ctx context.Context) (int64, error)
	PopAndReserve(ctx context.Context) (int64, error)
	Ack(ctx context.Context, taskID int64) error
	MoveToDead(ctx context.Context, taskID int64) error
}

type TaskRepository

type TaskRepository interface {
	Create(ctx context.Context, task *domain.Task) error
	GetByID(ctx context.Context, id int64) (*domain.Task, error)
	Update(ctx context.Context, task *domain.Task) error

	ListByParent(ctx context.Context, parentID int64) ([]*domain.Task, error)
	FindRunningRootTasks(ctx context.Context, before time.Time) ([]*domain.Task, error)

	FindByWorkflowID(ctx context.Context, workflowID int64) ([]*domain.Task, error)
	ListChildrenByParentID(ctx context.Context, parentID int64) ([]*domain.Task, error)

	// 批量更新更高效
	BatchUpdateStatus(ctx context.Context, taskIDs []int64, status domain.TaskStatus, errMsg string) error

	Enqueue(ctx context.Context, taskID int64) error
	// TryClaimTask 原子抢占任务(CAS)
	// 只允许一个 Worker 抢到任务
	// 同时支持 Running 超时任务重新抢占(Lease)
	TryClaimTask(ctx context.Context, taskID int64, workerID string) (bool, error)

	FindBySubKey(ctx context.Context, subKey string) (*domain.Task, error)
	ListByParentNode(ctx context.Context, parentID int64, nodeName string) ([]*domain.Task, error)

	CreateFork(ctx context.Context, source *domain.Task, newTaskID int64, newInput []byte, editAction, editLabel string) (*domain.Task, error)

	// 发布详情仍然取完整 task,但要求只能取 root task。返回 domain 类型,无 dto 依赖。
	GetRootTaskByIDAndUser(
		ctx context.Context,
		taskID int64,
		userID int64,
	) (*domain.Task, error)
}

TaskRepository 是引擎运行时依赖的任务存储接口,只涉及 domain 类型, 不引入任何面向 HTTP/API 的展示层(dto)依赖。面向业务的分页/详情查询 见 repository/query.TaskQueryRepository。

type WorkflowRepository

type WorkflowRepository interface {
	Create(ctx context.Context, wf *domain.Workflow) error
	Update(ctx context.Context, wf *domain.Workflow) error
	GetByID(ctx context.Context, id int64) (*domain.Workflow, error)
	GetByName(ctx context.Context, name string) (*domain.Workflow, error)
	List(ctx context.Context) ([]*domain.Workflow, error)
}

WorkflowRepository 存 工作流定义 元信息

type WorkflowVersionRepository

type WorkflowVersionRepository interface {
	Create(ctx context.Context, version *domain.WorkflowVersion) error
	Get(ctx context.Context, id int64) (*domain.WorkflowVersion, error)
	GetLatestByWorkflowID(ctx context.Context, id int64) (*domain.WorkflowVersion, error)
	GetLatestByWorkflowName(ctx context.Context, name string) (*domain.WorkflowVersion, error)
	UpdateDefinitionJSON(ctx context.Context, versionID int64, json []byte) error
}

Directories

Path Synopsis
redisqueue
Package redisqueue provides a Redis backed repository.TaskQueue.
Package redisqueue provides a Redis backed repository.TaskQueue.
taskapi
Package taskapi hosts the business/HTTP-facing task read-model queries (pagination, list, detail) that return presentation dto types.
Package taskapi hosts the business/HTTP-facing task read-model queries (pagination, list, detail) that return presentation dto types.

Jump to

Keyboard shortcuts

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