dispatcher

package
v0.110.13 Latest Latest
Warning

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

Go to latest
Published: Oct 7, 2026 License: MIT Imports: 48 Imported by: 0

Documentation

Index

Constants

View Source
const HeartbeatInterval = 4 * time.Second

Variables

View Source
var ErrNoActiveDurableInvocation = errors.New("no active durable invocation found")

ErrNoActiveDurableInvocation is returned by DeliverDurableEventLogEntryCompletion when this dispatcher has no in-memory durable session for the task -- typically because an engine restart wiped it and the worker has not reconnected yet, catching this error so that messages get sent do the DLQ lets DispatchCallbacks re-route it once the worker reconnects to a live dispatcher.

View Source
var ErrWorkerNotFound = fmt.Errorf("worker not found")

Functions

func UnmarshalPayload

func UnmarshalPayload[T any](payload interface{}) (T, error)

Types

type Dispatcher

type Dispatcher interface {
	contractsconnect.DispatcherHandler
	Start() (func() error, error)
}

type DispatcherImpl

type DispatcherImpl struct {
	contractsconnect.UnimplementedDispatcherHandler
	// contains filtered or unexported fields
}

func New

func New(fs ...DispatcherOpt) (*DispatcherImpl, error)

func (*DispatcherImpl) AddOperatorSession added in v0.109.0

func (d *DispatcherImpl) AddOperatorSession(
	workerId uuid.UUID,
	sessionId uuid.UUID,
	handler operator.ActionHandler,
) *OperatorHandlerSession

AddOperatorSession registers an in-process operator as a live session for workerId, so the dispatcher routes assigned actions to it like any other worker. The caller chooses sessionId so the dispatcher's session key is the same id it records on the worker row as the listener session fence. The caller must call Release when the session ends.

func (*DispatcherImpl) AddOperatorStreamSession added in v0.109.0

func (d *DispatcherImpl) AddOperatorStreamSession(
	ctx context.Context,
	workerId uuid.UUID,
	sessionId uuid.UUID,
	sender *rpcstream.Sender[v1contracts.OperatorListenResponse],
) *OperatorStreamSession

AddOperatorStreamSession registers a Listen stream owned by an out-of-process operator as a live session for workerId, so the dispatcher fans assigned actions out to it like any SDK worker. sender is the stream's guarded sender, which the handler closes before it returns. The caller chooses sessionId so the dispatcher's session key is the same id it records on the worker row as the listener session fence.

The dispatcher signals Fin when it wants the stream hung up (the shutdown drain in Start's cleanup). The caller must call Release when its handler exits; Release removes the session and lets a concurrent drain skip it, so a handler that has stopped selecting on Fin never blocks the drain. A send that is in progress when Release runs finishes or fails on its own once the handler closes the sender; Release does not wait for it.

func (*DispatcherImpl) CancelDAGChildren added in v0.104.0

func (d *DispatcherImpl) CancelDAGChildren(ctx context.Context, tenantId uuid.UUID, taskExternalIds []uuid.UUID) error

func (*DispatcherImpl) CancelStreamSessions added in v0.89.2

func (d *DispatcherImpl) CancelStreamSessions()

CancelStreamSessions hangs up all registered long-lived subscriber streams. It is called during shutdown before GracefulStop, which would otherwise block on them until the process is killed.

func (*DispatcherImpl) CancelTaskWithReason added in v0.109.0

func (d *DispatcherImpl) CancelTaskWithReason(ctx context.Context, tenantId uuid.UUID, request *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)

CancelTaskWithReason reports a cancelled task with a custom cancellation reason on behalf of an engine-internal operator (operator.TaskEventWriter). It is not a gRPC handler: the tenant is an argument, not the one the auth middleware puts on a request context.

func (*DispatcherImpl) DispatcherId added in v0.80.0

func (d *DispatcherImpl) DispatcherId() uuid.UUID

func (*DispatcherImpl) GetLocalWorkerIds added in v0.78.0

func (d *DispatcherImpl) GetLocalWorkerIds() map[uuid.UUID]struct{}

func (*DispatcherImpl) GetVersion added in v0.78.27

func (*DispatcherImpl) HandleLocalAssignments added in v0.78.0

func (d *DispatcherImpl) HandleLocalAssignments(ctx context.Context, tenantId, workerId uuid.UUID, tasks []*schedulingv1.AssignedItemWithTask) error

Note: this is very similar to handleTaskBulkAssignedTask, with some differences in what's sync vs run in a goroutine In this method, we wait until all tasks have been sent to the worker before returning

func (*DispatcherImpl) Heartbeat

Heartbeat is used to update the last heartbeat time for a worker

func (*DispatcherImpl) Listen

Subscribe handles a subscribe request from a client

func (*DispatcherImpl) ListenV2

ListenV2 is like Listen, but implementation does not include heartbeats. This should only used by SDKs against engine version v0.18.1+

func (*DispatcherImpl) NotifyNewWorker added in v0.106.4

func (d *DispatcherImpl) NotifyNewWorker(ctx context.Context, tenant *sqlcv1.Tenant, workerId uuid.UUID)

NotifyNewWorker tells the tenant's scheduler partition that workerId is available for work. The publish runs detached from ctx so it outlives the calling handler, and is a no-op for tenants without a scheduler partition.

func (*DispatcherImpl) PutOverridesData

func (*DispatcherImpl) RefreshTimeout

func (*DispatcherImpl) Register

func (*DispatcherImpl) RegisterDurableTask added in v0.104.0

func (d *DispatcherImpl) RegisterDurableTask(ctx context.Context, externalId uuid.UUID) (chan<- *contracts.DurableTaskRequest, <-chan *contracts.DurableTaskResponse, error)

RegisterDurableTask delegates to the V1 dispatcher service so DispatcherImpl satisfies operator.TaskEventWriter, giving operators the in-engine equivalent of the DurableTask RPC.

func (*DispatcherImpl) ReleaseSlot

func (*DispatcherImpl) RestoreEvictedTask added in v0.80.0

func (*DispatcherImpl) SendBatchActionEvent added in v0.98.0

func (s *DispatcherImpl) SendBatchActionEvent(ctx context.Context, request *contracts.BatchActionEvent) (*contracts.ActionEventResponse, error)

func (*DispatcherImpl) SendGroupKeyActionEvent

func (s *DispatcherImpl) SendGroupKeyActionEvent(ctx context.Context, request *contracts.GroupKeyActionEvent) (*contracts.ActionEventResponse, error)

func (*DispatcherImpl) SendStepActionEvent

func (s *DispatcherImpl) SendStepActionEvent(ctx context.Context, request *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)

func (*DispatcherImpl) Start

func (d *DispatcherImpl) Start() (func() error, error)

func (*DispatcherImpl) SubscribeToWorkflowEvents

func (s *DispatcherImpl) SubscribeToWorkflowEvents(ctx context.Context, request *contracts.SubscribeToWorkflowEventsRequest, connectStream *connect.ServerStream[contracts.WorkflowEvent]) error

func (*DispatcherImpl) SubscribeToWorkflowRuns

func (*DispatcherImpl) TriggerDAGStep added in v0.104.0

func (*DispatcherImpl) UpsertWorkerLabels

func (*DispatcherImpl) V1 added in v0.89.3

V1 returns the dispatcher's V1Dispatcher gRPC service, which serves the durable task and durable event RPCs.

type DispatcherOpt

type DispatcherOpt func(*DispatcherOpts)

func WithAlerter

func WithAlerter(a hatcheterrors.Alerter) DispatcherOpt

func WithAnalytics added in v0.79.36

func WithAnalytics(a analytics.Analytics) DispatcherOpt

func WithCache

func WithCache(cache cache.Cacheable) DispatcherOpt

func WithDefaultMaxWorkerLockAcquisitionTime added in v0.80.5

func WithDefaultMaxWorkerLockAcquisitionTime(t time.Duration) DispatcherOpt

func WithDispatcherId

func WithDispatcherId(dispatcherId uuid.UUID) DispatcherOpt

func WithLogger

func WithLogger(l *zerolog.Logger) DispatcherOpt

func WithMessageQueueV1

func WithMessageQueueV1(mqv1 msgqueue.MessageQueue) DispatcherOpt

func WithPayloadSizeThreshold

func WithPayloadSizeThreshold(threshold int) DispatcherOpt

func WithPrometheusGate added in v0.90.0

func WithPrometheusGate(gate *prometheus.Gate) DispatcherOpt

func WithPubSub added in v0.98.9

func WithPubSub(pubsub msgqueue.PubSub) DispatcherOpt

func WithRepositoryV1

func WithRepositoryV1(r v1.Repository) DispatcherOpt

func WithStreamEventBufferTimeout added in v0.79.34

func WithStreamEventBufferTimeout(timeout time.Duration) DispatcherOpt

func WithVersion added in v0.78.27

func WithVersion(version string) DispatcherOpt

func WithWorkflowRunBufferSize added in v0.78.0

func WithWorkflowRunBufferSize(size int) DispatcherOpt

type DispatcherOpts

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

type DispatcherServiceImpl added in v0.89.3

type DispatcherServiceImpl struct {
	v1connect.UnimplementedV1DispatcherHandler
	// contains filtered or unexported fields
}

func (*DispatcherServiceImpl) CancelDAGChildren added in v0.104.0

func (d *DispatcherServiceImpl) CancelDAGChildren(ctx context.Context, tenantId uuid.UUID, taskExternalIds []uuid.UUID) error

CancelDAGChildren cancels already-triggered DAG children on the operator's behalf, since it has no direct message-queue access; mirrors the admin CancelTasks path.

func (*DispatcherServiceImpl) CancelStreamSessions added in v0.89.3

func (d *DispatcherServiceImpl) CancelStreamSessions()

CancelStreamSessions hangs up all registered long-lived streams (durable event and durable task listeners). It is called during shutdown before GracefulStop, which would otherwise block on them until the process is killed.

func (*DispatcherServiceImpl) DeliverDurableEventLogEntryCompletion added in v0.89.3

func (d *DispatcherServiceImpl) DeliverDurableEventLogEntryCompletion(tenantId uuid.UUID, taskExternalId uuid.UUID, invocationCount int32, branchId, nodeId int64, payload []byte, satisfiedOrder *int64, isFailure bool, errorMessage *string) error

func (*DispatcherServiceImpl) DurableTask added in v0.89.3

func (*DispatcherServiceImpl) DurableTaskWithReceive added in v0.109.0

func (d *DispatcherServiceImpl) DurableTaskWithReceive(
	ctx context.Context,
	receive func() (*contracts.DurableTaskRequest, error),
	sender *rpcstream.Sender[contracts.DurableTaskResponse],
) error

DurableTaskWithReceive runs a durable task session over the two halves of a DurableTask stream that another service owns: the operator service authorizes the stream and reads its first message itself, then hands the rest over through receive. sender must be closed by the caller before its handler returns.

func (*DispatcherServiceImpl) ListenForDurableEvent added in v0.89.3

func (*DispatcherServiceImpl) RegisterDurableEvent added in v0.89.3

func (*DispatcherServiceImpl) RegisterDurableTask added in v0.104.0

func (d *DispatcherServiceImpl) RegisterDurableTask(ctx context.Context, externalId uuid.UUID) (chan<- *contracts.DurableTaskRequest, <-chan *contracts.DurableTaskResponse, error)

RegisterDurableTask sets up a channel-backed durable-task session — the in-engine equivalent of the DurableTask gRPC stream, for operators that don't hold a gRPC stream. The caller writes DurableTaskRequests to the returned requestCh and reads DurableTaskResponses from respCh, reusing the same handlers and routing table (durableInvocations) as the gRPC path, so async responses (wait-for satisfied, evictions, acks) are delivered identically.

externalId is registered up front so responses route immediately; additional task ids are registered lazily as messages reference them, matching DurableTask. The session is torn down (invocations deregistered, respCh closed) when ctx is cancelled or requestCh is closed.

func (*DispatcherServiceImpl) TriggerDAGStep added in v0.104.0

type OperatorHandlerSession added in v0.109.0

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

OperatorHandlerSession is a live dispatcher session backed by an operator running in this process: assigned actions are handed to its handler directly, with no encoding and no stream. There is no Fin signal, because there is no stream to hang up, and the shutdown drain skips these sessions: the host that opened the session owns its teardown.

func (*OperatorHandlerSession) Release added in v0.109.0

func (s *OperatorHandlerSession) Release()

Release removes the session from the dispatcher. It is idempotent.

func (*OperatorHandlerSession) SetPaused added in v0.109.0

func (s *OperatorHandlerSession) SetPaused(paused bool)

SetPaused is the same as OperatorStreamSession.SetPaused: starts assigned to the worker are returned to the queue instead of handed to the handler.

type OperatorStreamSession added in v0.109.0

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

OperatorStreamSession is a live dispatcher session backed by an OperatorService Listen stream that an operator service owns. The dispatcher fans assigned actions out on the stream; the owner sends its own protocol messages through Send. Both go through the one sender that guards the stream, hangs up when Fin fires, and pauses delivery through SetPaused.

func (*OperatorStreamSession) Fin added in v0.109.0

func (s *OperatorStreamSession) Fin() <-chan bool

Fin fires when the dispatcher wants the stream hung up.

func (*OperatorStreamSession) Release added in v0.109.0

func (s *OperatorStreamSession) Release()

Release removes the session from the dispatcher. It is idempotent.

func (*OperatorStreamSession) Send added in v0.109.0

Send writes msg on the stream, serialised with the dispatcher's own action sends. It fails with errSessionReleased after Release and with rpcstream.ErrClosed once the handler has closed the sender.

func (*OperatorStreamSession) SetPaused added in v0.109.0

func (s *OperatorStreamSession) SetPaused(paused bool)

SetPaused makes the dispatcher return every start assigned to the worker to the queue instead of sending it, until SetPaused(false). The owner sets it before it acknowledges a pause to the operator, so an action the scheduler assigned before it observed the pause is requeued rather than delivered after the ack. Cancels are still sent.

type StreamEventBuffer

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

func NewStreamEventBuffer

func NewStreamEventBuffer(timeout time.Duration) *StreamEventBuffer

func (*StreamEventBuffer) AddEvent

func (b *StreamEventBuffer) AddEvent(event *contracts.WorkflowEvent)

func (*StreamEventBuffer) Close

func (b *StreamEventBuffer) Close()

Close stops the buffer's background goroutines. The channels are deliberately not closed: in-flight AddEvent/timeout sends race a close (send on a closed channel panics even inside a select), and every sender and consumer already exits via context cancellation.

func (*StreamEventBuffer) Events

func (b *StreamEventBuffer) Events() <-chan *contracts.WorkflowEvent

type V1TaskWithPayloadAndInvocationCount added in v0.80.0

type V1TaskWithPayloadAndInvocationCount struct {
	*v1.V1TaskWithPayload
	InvocationCount *int32 // only used for durable tasks
}

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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