Documentation
¶
Index ¶
- Variables
- type TriggerWriter
- func (tw *TriggerWriter) SignalCELEvaluationFailures(ctx context.Context, tenantId uuid.UUID, failures []v1.CELEvaluationFailure) error
- func (tw *TriggerWriter) SignalCreated(ctx context.Context, tenantId uuid.UUID, tasks []*v1.V1TaskWithPayload, ...) error
- func (tw *TriggerWriter) TriggerFromEvents(ctx context.Context, tenantId uuid.UUID, ...) error
- func (tw *TriggerWriter) TriggerFromWorkflowNames(ctx context.Context, tenantId uuid.UUID, opts []*v1.WorkflowNameTriggerOpts) ([]v1.IdempotencyCollision, error)
- func (tw *TriggerWriter) TriggerFromWorkflowNamesWaiting(ctx context.Context, tenantId uuid.UUID, opts []*v1.WorkflowNameTriggerOpts) ([]v1.IdempotencyCollision, error)
Constants ¶
This section is empty.
Variables ¶
View Source
var ErrNoTriggerSlots = errors.New("no trigger slots available")
Functions ¶
This section is empty.
Types ¶
type TriggerWriter ¶
type TriggerWriter struct {
// contains filtered or unexported fields
}
func NewTriggerWriter ¶
func NewTriggerWriter(mq msgqueue.MessageQueue, pubsub msgqueue.PubSub, repo v1.Repository, l *zerolog.Logger, pubBuffer *msgqueue.MQPubBuffer, slots int, promGate *prometheus.Gate) *TriggerWriter
NewTriggerWriter creates a new TriggerWriter with the given number of slots for concurrency control. If the number of slots is 0, there is no limit to concurrency.
func (*TriggerWriter) SignalCELEvaluationFailures ¶ added in v0.95.0
func (tw *TriggerWriter) SignalCELEvaluationFailures(ctx context.Context, tenantId uuid.UUID, failures []v1.CELEvaluationFailure) error
func (*TriggerWriter) SignalCreated ¶ added in v0.80.0
func (tw *TriggerWriter) SignalCreated(ctx context.Context, tenantId uuid.UUID, tasks []*v1.V1TaskWithPayload, dags []*v1.DAGWithData) error
func (*TriggerWriter) TriggerFromEvents ¶
func (tw *TriggerWriter) TriggerFromEvents(ctx context.Context, tenantId uuid.UUID, eventIdToOpts map[uuid.UUID]v1.EventTriggerOpts) error
func (*TriggerWriter) TriggerFromWorkflowNames ¶
func (tw *TriggerWriter) TriggerFromWorkflowNames(ctx context.Context, tenantId uuid.UUID, opts []*v1.WorkflowNameTriggerOpts) ([]v1.IdempotencyCollision, error)
func (*TriggerWriter) TriggerFromWorkflowNamesWaiting ¶ added in v0.108.2
func (tw *TriggerWriter) TriggerFromWorkflowNamesWaiting(ctx context.Context, tenantId uuid.UUID, opts []*v1.WorkflowNameTriggerOpts) ([]v1.IdempotencyCollision, error)
TriggerFromWorkflowNamesWaiting acquires a trigger slot, blocking until one is free or ctx is done. Use this when the caller cannot fall back to the durable queue (idempotent batches).
Click to show internal directories.
Click to hide internal directories.