Documentation
¶
Index ¶
- Variables
- type ProcGroup
- func (pg *ProcGroup) Acquire(_ context.Context, dagRun execution.DAGRunRef) (*ProcHandle, error)
- func (pg *ProcGroup) Count(ctx context.Context) (int, error)
- func (pg *ProcGroup) CountByDAGName(ctx context.Context, dagName string) (int, error)
- func (pg *ProcGroup) IsRunAlive(ctx context.Context, dagRun execution.DAGRunRef) (bool, error)
- func (pg *ProcGroup) ListAlive(ctx context.Context) ([]execution.DAGRunRef, error)
- type ProcHandle
- type Store
- func (s *Store) Acquire(ctx context.Context, groupName string, dagRun execution.DAGRunRef) (execution.ProcHandle, error)
- func (s *Store) CountAlive(ctx context.Context, groupName string) (int, error)
- func (s *Store) CountAliveByDAGName(ctx context.Context, groupName, dagName string) (int, error)
- func (s *Store) IsRunAlive(ctx context.Context, groupName string, dagRun execution.DAGRunRef) (bool, error)
- func (s *Store) ListAlive(ctx context.Context, groupName string) ([]execution.DAGRunRef, error)
- func (s *Store) ListAllAlive(ctx context.Context) (map[string][]execution.DAGRunRef, error)
- func (s *Store) Lock(ctx context.Context, groupName string) error
- func (s *Store) Unlock(ctx context.Context, groupName string)
Constants ¶
This section is empty.
Variables ¶
var (
ErrHeartbeatAlreadyStarted = fmt.Errorf("heartbeat already started")
)
Error messages
Functions ¶
This section is empty.
Types ¶
type ProcGroup ¶
ProcGroup is a struct that manages process files for a given DAG name.
func NewProcGroup ¶
NewProcGroup creates a new instance of a ProcGroup with the specified base directory and DAG name.
func (*ProcGroup) Acquire ¶
GetProc retrieves a proc file for the specified dag-run reference. It returns a new Proc instance with the generated file name.
func (*ProcGroup) CountByDAGName ¶ added in v1.22.0
func (*ProcGroup) IsRunAlive ¶ added in v1.18.0
IsRunAlive checks if a specific DAG run has an alive process file.
type ProcHandle ¶
type ProcHandle struct {
// contains filtered or unexported fields
}
ProcHandle is a struct that implements the ProcHandle interface.
func NewProcHandler ¶
func NewProcHandler(file string, meta execution.ProcMeta) *ProcHandle
NewProcHandler creates a new instance of Proc with the specified file name.
func (*ProcHandle) GetMeta ¶
func (p *ProcHandle) GetMeta() execution.ProcMeta
GetMeta implements models.ProcHandle.
type Store ¶
type Store struct {
// contains filtered or unexported fields
}
Store is a struct that implements the ProcStore interface.
func (*Store) Acquire ¶
func (s *Store) Acquire(ctx context.Context, groupName string, dagRun execution.DAGRunRef) (execution.ProcHandle, error)
Acquire implements models.ProcStore.
func (*Store) CountAlive ¶
CountAlive implements models.ProcStore.
func (*Store) CountAliveByDAGName ¶ added in v1.22.0
func (*Store) IsRunAlive ¶ added in v1.18.0
func (s *Store) IsRunAlive(ctx context.Context, groupName string, dagRun execution.DAGRunRef) (bool, error)
IsRunAlive implements models.ProcStore.
func (*Store) ListAllAlive ¶ added in v1.22.0
ListAllAlive implements models.ProcStore. Returns all running DAG runs across all process groups.