Documentation
¶
Index ¶
- Constants
- Variables
- func UnmarshalPayload[T any](payload interface{}) (T, error)
- type Dispatcher
- type DispatcherImpl
- func (d *DispatcherImpl) AddOperatorSession(workerId uuid.UUID, sessionId uuid.UUID, handler operator.ActionHandler) *OperatorHandlerSession
- func (d *DispatcherImpl) AddOperatorStreamSession(ctx context.Context, workerId uuid.UUID, sessionId uuid.UUID, ...) *OperatorStreamSession
- func (d *DispatcherImpl) CancelDAGChildren(ctx context.Context, tenantId uuid.UUID, taskExternalIds []uuid.UUID) error
- func (d *DispatcherImpl) CancelStreamSessions()
- func (d *DispatcherImpl) CancelTaskWithReason(ctx context.Context, tenantId uuid.UUID, request *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)
- func (d *DispatcherImpl) DispatcherId() uuid.UUID
- func (d *DispatcherImpl) GetLocalWorkerIds() map[uuid.UUID]struct{}
- func (s *DispatcherImpl) GetVersion(ctx context.Context, req *contracts.GetVersionRequest) (*contracts.GetVersionResponse, error)
- func (d *DispatcherImpl) HandleLocalAssignments(ctx context.Context, tenantId, workerId uuid.UUID, ...) error
- func (s *DispatcherImpl) Heartbeat(ctx context.Context, req *contracts.HeartbeatRequest) (*contracts.HeartbeatResponse, error)
- func (s *DispatcherImpl) Listen(ctx context.Context, request *contracts.WorkerListenRequest, ...) error
- func (s *DispatcherImpl) ListenV2(ctx context.Context, request *contracts.WorkerListenRequest, ...) error
- func (d *DispatcherImpl) NotifyNewWorker(ctx context.Context, tenant *sqlcv1.Tenant, workerId uuid.UUID)
- func (s *DispatcherImpl) PutOverridesData(ctx context.Context, request *contracts.OverridesData) (*contracts.OverridesDataResponse, error)
- func (d *DispatcherImpl) RefreshTimeout(ctx context.Context, request *contracts.RefreshTimeoutRequest) (*contracts.RefreshTimeoutResponse, error)
- func (s *DispatcherImpl) Register(ctx context.Context, request *contracts.WorkerRegisterRequest) (*contracts.WorkerRegisterResponse, error)
- func (d *DispatcherImpl) RegisterDurableTask(ctx context.Context, externalId uuid.UUID) (chan<- *contracts.DurableTaskRequest, <-chan *contracts.DurableTaskResponse, ...)
- func (s *DispatcherImpl) ReleaseSlot(ctx context.Context, req *contracts.ReleaseSlotRequest) (*contracts.ReleaseSlotResponse, error)
- func (s *DispatcherImpl) RestoreEvictedTask(ctx context.Context, req *contracts.RestoreEvictedTaskRequest) (*contracts.RestoreEvictedTaskResponse, error)
- func (s *DispatcherImpl) SendBatchActionEvent(ctx context.Context, request *contracts.BatchActionEvent) (*contracts.ActionEventResponse, error)
- func (s *DispatcherImpl) SendGroupKeyActionEvent(ctx context.Context, request *contracts.GroupKeyActionEvent) (*contracts.ActionEventResponse, error)
- func (s *DispatcherImpl) SendStepActionEvent(ctx context.Context, request *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)
- func (d *DispatcherImpl) Start() (func() error, error)
- func (s *DispatcherImpl) SubscribeToWorkflowEvents(ctx context.Context, request *contracts.SubscribeToWorkflowEventsRequest, ...) error
- func (s *DispatcherImpl) SubscribeToWorkflowRuns(ctx context.Context, ...) error
- func (d *DispatcherImpl) TriggerDAGStep(ctx context.Context, tenantId uuid.UUID, req *operator.DAGStepTriggerRequest) (*operator.DAGStepTriggerResult, error)
- func (s *DispatcherImpl) Unsubscribe(ctx context.Context, request *contracts.WorkerUnsubscribeRequest) (*contracts.WorkerUnsubscribeResponse, error)
- func (s *DispatcherImpl) UpsertWorkerLabels(ctx context.Context, request *contracts.UpsertWorkerLabelsRequest) (*contracts.UpsertWorkerLabelsResponse, error)
- func (d *DispatcherImpl) V1() *DispatcherServiceImpl
- type DispatcherOpt
- func WithAlerter(a hatcheterrors.Alerter) DispatcherOpt
- func WithAnalytics(a analytics.Analytics) DispatcherOpt
- func WithCache(cache cache.Cacheable) DispatcherOpt
- func WithDataDecoderValidator(dv datautils.DataDecoderValidator) DispatcherOpt
- func WithDefaultMaxWorkerLockAcquisitionTime(t time.Duration) DispatcherOpt
- func WithDispatcherId(dispatcherId uuid.UUID) DispatcherOpt
- func WithLogger(l *zerolog.Logger) DispatcherOpt
- func WithMessageQueueV1(mqv1 msgqueue.MessageQueue) DispatcherOpt
- func WithPayloadSizeThreshold(threshold int) DispatcherOpt
- func WithPrometheusGate(gate *prometheus.Gate) DispatcherOpt
- func WithPubSub(pubsub msgqueue.PubSub) DispatcherOpt
- func WithRepositoryV1(r v1.Repository) DispatcherOpt
- func WithStreamEventBufferTimeout(timeout time.Duration) DispatcherOpt
- func WithVersion(version string) DispatcherOpt
- func WithWorkflowRunBufferSize(size int) DispatcherOpt
- type DispatcherOpts
- type DispatcherServiceImpl
- func (d *DispatcherServiceImpl) CancelDAGChildren(ctx context.Context, tenantId uuid.UUID, taskExternalIds []uuid.UUID) error
- func (d *DispatcherServiceImpl) CancelStreamSessions()
- func (d *DispatcherServiceImpl) DeliverDurableEventLogEntryCompletion(tenantId uuid.UUID, taskExternalId uuid.UUID, invocationCount int32, ...) error
- func (d *DispatcherServiceImpl) DurableTask(ctx context.Context, ...) error
- func (d *DispatcherServiceImpl) DurableTaskWithReceive(ctx context.Context, receive func() (*contracts.DurableTaskRequest, error), ...) error
- func (d *DispatcherServiceImpl) ListenForDurableEvent(ctx context.Context, ...) error
- func (d *DispatcherServiceImpl) RegisterDurableEvent(ctx context.Context, req *contracts.RegisterDurableEventRequest) (*contracts.RegisterDurableEventResponse, error)
- func (d *DispatcherServiceImpl) RegisterDurableTask(ctx context.Context, externalId uuid.UUID) (chan<- *contracts.DurableTaskRequest, <-chan *contracts.DurableTaskResponse, ...)
- func (d *DispatcherServiceImpl) TriggerDAGStep(ctx context.Context, tenantId uuid.UUID, req *operator.DAGStepTriggerRequest) (*operator.DAGStepTriggerResult, error)
- type OperatorHandlerSession
- type OperatorStreamSession
- type StreamEventBuffer
- type V1TaskWithPayloadAndInvocationCount
Constants ¶
const HeartbeatInterval = 4 * time.Second
Variables ¶
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.
var ErrWorkerNotFound = fmt.Errorf("worker not found")
Functions ¶
func UnmarshalPayload ¶
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 (*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 (s *DispatcherImpl) GetVersion(ctx context.Context, req *contracts.GetVersionRequest) (*contracts.GetVersionResponse, error)
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 ¶
func (s *DispatcherImpl) Heartbeat(ctx context.Context, req *contracts.HeartbeatRequest) (*contracts.HeartbeatResponse, error)
Heartbeat is used to update the last heartbeat time for a worker
func (*DispatcherImpl) Listen ¶
func (s *DispatcherImpl) Listen(ctx context.Context, request *contracts.WorkerListenRequest, connectStream *connect.ServerStream[contracts.AssignedAction]) error
Subscribe handles a subscribe request from a client
func (*DispatcherImpl) ListenV2 ¶
func (s *DispatcherImpl) ListenV2(ctx context.Context, request *contracts.WorkerListenRequest, connectStream *connect.ServerStream[contracts.AssignedAction]) error
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 (s *DispatcherImpl) PutOverridesData(ctx context.Context, request *contracts.OverridesData) (*contracts.OverridesDataResponse, error)
func (*DispatcherImpl) RefreshTimeout ¶
func (d *DispatcherImpl) RefreshTimeout(ctx context.Context, request *contracts.RefreshTimeoutRequest) (*contracts.RefreshTimeoutResponse, error)
func (*DispatcherImpl) Register ¶
func (s *DispatcherImpl) Register(ctx context.Context, request *contracts.WorkerRegisterRequest) (*contracts.WorkerRegisterResponse, error)
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 (s *DispatcherImpl) ReleaseSlot(ctx context.Context, req *contracts.ReleaseSlotRequest) (*contracts.ReleaseSlotResponse, error)
func (*DispatcherImpl) RestoreEvictedTask ¶ added in v0.80.0
func (s *DispatcherImpl) RestoreEvictedTask(ctx context.Context, req *contracts.RestoreEvictedTaskRequest) (*contracts.RestoreEvictedTaskResponse, error)
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 (s *DispatcherImpl) SubscribeToWorkflowRuns(ctx context.Context, connectStream *connect.BidiStream[contracts.SubscribeToWorkflowRunsRequest, contracts.WorkflowRunEvent]) error
func (*DispatcherImpl) TriggerDAGStep ¶ added in v0.104.0
func (d *DispatcherImpl) TriggerDAGStep(ctx context.Context, tenantId uuid.UUID, req *operator.DAGStepTriggerRequest) (*operator.DAGStepTriggerResult, error)
func (*DispatcherImpl) Unsubscribe ¶
func (s *DispatcherImpl) Unsubscribe(ctx context.Context, request *contracts.WorkerUnsubscribeRequest) (*contracts.WorkerUnsubscribeResponse, error)
func (*DispatcherImpl) UpsertWorkerLabels ¶
func (s *DispatcherImpl) UpsertWorkerLabels(ctx context.Context, request *contracts.UpsertWorkerLabelsRequest) (*contracts.UpsertWorkerLabelsResponse, error)
func (*DispatcherImpl) V1 ¶ added in v0.89.3
func (d *DispatcherImpl) V1() *DispatcherServiceImpl
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 WithDataDecoderValidator ¶
func WithDataDecoderValidator(dv datautils.DataDecoderValidator) 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 (*DispatcherServiceImpl) DurableTask ¶ added in v0.89.3
func (d *DispatcherServiceImpl) DurableTask(ctx context.Context, server *connect.BidiStream[contracts.DurableTaskRequest, contracts.DurableTaskResponse]) error
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 (d *DispatcherServiceImpl) ListenForDurableEvent(ctx context.Context, server *connect.BidiStream[contracts.ListenForDurableEventRequest, contracts.DurableEvent]) error
func (*DispatcherServiceImpl) RegisterDurableEvent ¶ added in v0.89.3
func (d *DispatcherServiceImpl) RegisterDurableEvent(ctx context.Context, req *contracts.RegisterDurableEventRequest) (*contracts.RegisterDurableEventResponse, error)
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
func (d *DispatcherServiceImpl) TriggerDAGStep(ctx context.Context, tenantId uuid.UUID, req *operator.DAGStepTriggerRequest) (*operator.DAGStepTriggerResult, error)
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
func (s *OperatorStreamSession) Send(ctx context.Context, msg *v1contracts.OperatorListenResponse) error
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
}