Documentation
¶
Index ¶
- Variables
- func RegisterAdminServiceServer(s grpc.ServiceRegistrar, srv AdminServiceServer)
- func RegisterOperatorServiceServer(s grpc.ServiceRegistrar, srv OperatorServiceServer)
- func RegisterV1DispatcherServer(s grpc.ServiceRegistrar, srv V1DispatcherServer)
- func RegisterV1StreamsServer(s grpc.ServiceRegistrar, srv V1StreamsServer)
- type Action
- type AdminServiceClient
- type AdminServiceServer
- type BaseMatchCondition
- func (*BaseMatchCondition) Descriptor() ([]byte, []int)deprecated
- func (x *BaseMatchCondition) GetAction() Action
- func (x *BaseMatchCondition) GetExpression() string
- func (x *BaseMatchCondition) GetOrGroupId() string
- func (x *BaseMatchCondition) GetReadableDataKey() string
- func (*BaseMatchCondition) ProtoMessage()
- func (x *BaseMatchCondition) ProtoReflect() protoreflect.Message
- func (x *BaseMatchCondition) Reset()
- func (x *BaseMatchCondition) String() string
- type BranchDurableTaskRequest
- func (*BranchDurableTaskRequest) Descriptor() ([]byte, []int)deprecated
- func (x *BranchDurableTaskRequest) GetBranchId() int64
- func (x *BranchDurableTaskRequest) GetNodeId() int64
- func (x *BranchDurableTaskRequest) GetTaskExternalId() string
- func (*BranchDurableTaskRequest) ProtoMessage()
- func (x *BranchDurableTaskRequest) ProtoReflect() protoreflect.Message
- func (x *BranchDurableTaskRequest) Reset()
- func (x *BranchDurableTaskRequest) String() string
- type BranchDurableTaskResponse
- func (*BranchDurableTaskResponse) Descriptor() ([]byte, []int)deprecated
- func (x *BranchDurableTaskResponse) GetBranchId() int64
- func (x *BranchDurableTaskResponse) GetNodeId() int64
- func (x *BranchDurableTaskResponse) GetTaskExternalId() string
- func (*BranchDurableTaskResponse) ProtoMessage()
- func (x *BranchDurableTaskResponse) ProtoReflect() protoreflect.Message
- func (x *BranchDurableTaskResponse) Reset()
- func (x *BranchDurableTaskResponse) String() string
- type BulkTriggerIdempotencyCollisionError
- func (*BulkTriggerIdempotencyCollisionError) Descriptor() ([]byte, []int)deprecated
- func (x *BulkTriggerIdempotencyCollisionError) GetCollisions() []*IdempotencyCollisionError
- func (x *BulkTriggerIdempotencyCollisionError) GetSuccessfulWorkflowRunExternalIds() []string
- func (*BulkTriggerIdempotencyCollisionError) ProtoMessage()
- func (x *BulkTriggerIdempotencyCollisionError) ProtoReflect() protoreflect.Message
- func (x *BulkTriggerIdempotencyCollisionError) Reset()
- func (x *BulkTriggerIdempotencyCollisionError) String() string
- type CancelTasksRequest
- func (*CancelTasksRequest) Descriptor() ([]byte, []int)deprecated
- func (x *CancelTasksRequest) GetExternalIds() []string
- func (x *CancelTasksRequest) GetFilter() *TasksFilter
- func (*CancelTasksRequest) ProtoMessage()
- func (x *CancelTasksRequest) ProtoReflect() protoreflect.Message
- func (x *CancelTasksRequest) Reset()
- func (x *CancelTasksRequest) String() string
- type CancelTasksResponse
- func (*CancelTasksResponse) Descriptor() ([]byte, []int)deprecated
- func (x *CancelTasksResponse) GetCancelledTasks() []string
- func (*CancelTasksResponse) ProtoMessage()
- func (x *CancelTasksResponse) ProtoReflect() protoreflect.Message
- func (x *CancelTasksResponse) Reset()
- func (x *CancelTasksResponse) String() string
- type Concurrency
- func (*Concurrency) Descriptor() ([]byte, []int)deprecated
- func (x *Concurrency) GetExpression() string
- func (x *Concurrency) GetIsTenantScoped() bool
- func (x *Concurrency) GetLimitStrategy() ConcurrencyLimitStrategy
- func (x *Concurrency) GetMaxRuns() int32
- func (x *Concurrency) GetMaxRunsExpression() string
- func (x *Concurrency) GetName() string
- func (*Concurrency) ProtoMessage()
- func (x *Concurrency) ProtoReflect() protoreflect.Message
- func (x *Concurrency) Reset()
- func (x *Concurrency) String() string
- type ConcurrencyLimitStrategy
- func (ConcurrencyLimitStrategy) Descriptor() protoreflect.EnumDescriptor
- func (x ConcurrencyLimitStrategy) Enum() *ConcurrencyLimitStrategy
- func (ConcurrencyLimitStrategy) EnumDescriptor() ([]byte, []int)deprecated
- func (x ConcurrencyLimitStrategy) Number() protoreflect.EnumNumber
- func (x ConcurrencyLimitStrategy) String() string
- func (ConcurrencyLimitStrategy) Type() protoreflect.EnumType
- type CreateTaskOpts
- func (*CreateTaskOpts) Descriptor() ([]byte, []int)deprecated
- func (x *CreateTaskOpts) GetAction() string
- func (x *CreateTaskOpts) GetBackoffFactor() float32
- func (x *CreateTaskOpts) GetBackoffMaxSeconds() int32
- func (x *CreateTaskOpts) GetBatch() *TaskBatchConfig
- func (x *CreateTaskOpts) GetConcurrency() []*Concurrency
- func (x *CreateTaskOpts) GetConditions() *TaskConditions
- func (x *CreateTaskOpts) GetInputs() string
- func (x *CreateTaskOpts) GetIsDurable() bool
- func (x *CreateTaskOpts) GetParents() []string
- func (x *CreateTaskOpts) GetRateLimits() []*CreateTaskRateLimit
- func (x *CreateTaskOpts) GetReadableId() string
- func (x *CreateTaskOpts) GetRetries() int32
- func (x *CreateTaskOpts) GetScheduleTimeout() string
- func (x *CreateTaskOpts) GetSlotRequests() map[string]int32
- func (x *CreateTaskOpts) GetTimeout() string
- func (x *CreateTaskOpts) GetWorkerLabels() map[string]*DesiredWorkerLabels
- func (*CreateTaskOpts) ProtoMessage()
- func (x *CreateTaskOpts) ProtoReflect() protoreflect.Message
- func (x *CreateTaskOpts) Reset()
- func (x *CreateTaskOpts) String() string
- type CreateTaskRateLimit
- func (*CreateTaskRateLimit) Descriptor() ([]byte, []int)deprecated
- func (x *CreateTaskRateLimit) GetDuration() RateLimitDuration
- func (x *CreateTaskRateLimit) GetKey() string
- func (x *CreateTaskRateLimit) GetKeyExpr() string
- func (x *CreateTaskRateLimit) GetLimitValuesExpr() string
- func (x *CreateTaskRateLimit) GetUnits() int32
- func (x *CreateTaskRateLimit) GetUnitsExpr() string
- func (*CreateTaskRateLimit) ProtoMessage()
- func (x *CreateTaskRateLimit) ProtoReflect() protoreflect.Message
- func (x *CreateTaskRateLimit) Reset()
- func (x *CreateTaskRateLimit) String() string
- type CreateWorkflowVersionRequest
- func (*CreateWorkflowVersionRequest) Descriptor() ([]byte, []int)deprecated
- func (x *CreateWorkflowVersionRequest) GetConcurrency() *Concurrency
- func (x *CreateWorkflowVersionRequest) GetConcurrencyArr() []*Concurrency
- func (x *CreateWorkflowVersionRequest) GetCronInput() string
- func (x *CreateWorkflowVersionRequest) GetCronTriggers() []string
- func (x *CreateWorkflowVersionRequest) GetDefaultFilters() []*DefaultFilter
- func (x *CreateWorkflowVersionRequest) GetDefaultPriority() int32
- func (x *CreateWorkflowVersionRequest) GetDescription() string
- func (x *CreateWorkflowVersionRequest) GetEventTriggers() []string
- func (x *CreateWorkflowVersionRequest) GetIdempotency() *IdempotencyConfig
- func (x *CreateWorkflowVersionRequest) GetInputJsonSchema() []byte
- func (x *CreateWorkflowVersionRequest) GetName() string
- func (x *CreateWorkflowVersionRequest) GetOnFailureTask() *CreateTaskOpts
- func (x *CreateWorkflowVersionRequest) GetSticky() StickyStrategy
- func (x *CreateWorkflowVersionRequest) GetTasks() []*CreateTaskOpts
- func (x *CreateWorkflowVersionRequest) GetVersion() string
- func (*CreateWorkflowVersionRequest) ProtoMessage()
- func (x *CreateWorkflowVersionRequest) ProtoReflect() protoreflect.Message
- func (x *CreateWorkflowVersionRequest) Reset()
- func (x *CreateWorkflowVersionRequest) String() string
- type CreateWorkflowVersionResponse
- func (*CreateWorkflowVersionResponse) Descriptor() ([]byte, []int)deprecated
- func (x *CreateWorkflowVersionResponse) GetId() string
- func (x *CreateWorkflowVersionResponse) GetWorkflowId() string
- func (*CreateWorkflowVersionResponse) ProtoMessage()
- func (x *CreateWorkflowVersionResponse) ProtoReflect() protoreflect.Message
- func (x *CreateWorkflowVersionResponse) Reset()
- func (x *CreateWorkflowVersionResponse) String() string
- type DefaultFilter
- func (*DefaultFilter) Descriptor() ([]byte, []int)deprecated
- func (x *DefaultFilter) GetExpression() string
- func (x *DefaultFilter) GetPayload() []byte
- func (x *DefaultFilter) GetScope() string
- func (*DefaultFilter) ProtoMessage()
- func (x *DefaultFilter) ProtoReflect() protoreflect.Message
- func (x *DefaultFilter) Reset()
- func (x *DefaultFilter) String() string
- type DesiredWorkerLabels
- func (*DesiredWorkerLabels) Descriptor() ([]byte, []int)deprecated
- func (x *DesiredWorkerLabels) GetComparator() WorkerLabelComparator
- func (x *DesiredWorkerLabels) GetIntValue() int32
- func (x *DesiredWorkerLabels) GetRequired() bool
- func (x *DesiredWorkerLabels) GetStrValue() string
- func (x *DesiredWorkerLabels) GetWeight() int32
- func (*DesiredWorkerLabels) ProtoMessage()
- func (x *DesiredWorkerLabels) ProtoReflect() protoreflect.Message
- func (x *DesiredWorkerLabels) Reset()
- func (x *DesiredWorkerLabels) String() string
- type DurableEvent
- func (*DurableEvent) Descriptor() ([]byte, []int)deprecated
- func (x *DurableEvent) GetData() []byte
- func (x *DurableEvent) GetSignalKey() string
- func (x *DurableEvent) GetTaskId() string
- func (*DurableEvent) ProtoMessage()
- func (x *DurableEvent) ProtoReflect() protoreflect.Message
- func (x *DurableEvent) Reset()
- func (x *DurableEvent) String() string
- type DurableEventListenerConditions
- func (*DurableEventListenerConditions) Descriptor() ([]byte, []int)deprecated
- func (x *DurableEventListenerConditions) GetSleepConditions() []*SleepMatchCondition
- func (x *DurableEventListenerConditions) GetUserEventConditions() []*UserEventMatchCondition
- func (*DurableEventListenerConditions) ProtoMessage()
- func (x *DurableEventListenerConditions) ProtoReflect() protoreflect.Message
- func (x *DurableEventListenerConditions) Reset()
- func (x *DurableEventListenerConditions) String() string
- type DurableEventLogEntryRef
- func (*DurableEventLogEntryRef) Descriptor() ([]byte, []int)deprecated
- func (x *DurableEventLogEntryRef) GetBranchId() int64
- func (x *DurableEventLogEntryRef) GetDurableTaskExternalId() string
- func (x *DurableEventLogEntryRef) GetInvocationCount() int32
- func (x *DurableEventLogEntryRef) GetNodeId() int64
- func (*DurableEventLogEntryRef) ProtoMessage()
- func (x *DurableEventLogEntryRef) ProtoReflect() protoreflect.Message
- func (x *DurableEventLogEntryRef) Reset()
- func (x *DurableEventLogEntryRef) String() string
- type DurableTaskAwaitedCompletedEntry
- func (*DurableTaskAwaitedCompletedEntry) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskAwaitedCompletedEntry) GetBranchId() int64
- func (x *DurableTaskAwaitedCompletedEntry) GetDurableTaskExternalId() string
- func (x *DurableTaskAwaitedCompletedEntry) GetInvocationCount() int32
- func (x *DurableTaskAwaitedCompletedEntry) GetNodeId() int64
- func (*DurableTaskAwaitedCompletedEntry) ProtoMessage()
- func (x *DurableTaskAwaitedCompletedEntry) ProtoReflect() protoreflect.Message
- func (x *DurableTaskAwaitedCompletedEntry) Reset()
- func (x *DurableTaskAwaitedCompletedEntry) String() string
- type DurableTaskCompleteMemoRequest
- func (*DurableTaskCompleteMemoRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskCompleteMemoRequest) GetMemoKey() []byte
- func (x *DurableTaskCompleteMemoRequest) GetPayload() []byte
- func (x *DurableTaskCompleteMemoRequest) GetRef() *DurableEventLogEntryRef
- func (*DurableTaskCompleteMemoRequest) ProtoMessage()
- func (x *DurableTaskCompleteMemoRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskCompleteMemoRequest) Reset()
- func (x *DurableTaskCompleteMemoRequest) String() string
- type DurableTaskErrorResponse
- func (*DurableTaskErrorResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskErrorResponse) GetErrorMessage() string
- func (x *DurableTaskErrorResponse) GetErrorType() DurableTaskErrorType
- func (x *DurableTaskErrorResponse) GetRef() *DurableEventLogEntryRef
- func (*DurableTaskErrorResponse) ProtoMessage()
- func (x *DurableTaskErrorResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskErrorResponse) Reset()
- func (x *DurableTaskErrorResponse) String() string
- type DurableTaskErrorType
- func (DurableTaskErrorType) Descriptor() protoreflect.EnumDescriptor
- func (x DurableTaskErrorType) Enum() *DurableTaskErrorType
- func (DurableTaskErrorType) EnumDescriptor() ([]byte, []int)deprecated
- func (x DurableTaskErrorType) Number() protoreflect.EnumNumber
- func (x DurableTaskErrorType) String() string
- func (DurableTaskErrorType) Type() protoreflect.EnumType
- type DurableTaskEventLogEntryCompletedResponse
- func (*DurableTaskEventLogEntryCompletedResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEventLogEntryCompletedResponse) GetErrorMessage() string
- func (x *DurableTaskEventLogEntryCompletedResponse) GetIsFailure() bool
- func (x *DurableTaskEventLogEntryCompletedResponse) GetPayload() []byte
- func (x *DurableTaskEventLogEntryCompletedResponse) GetRef() *DurableEventLogEntryRef
- func (*DurableTaskEventLogEntryCompletedResponse) ProtoMessage()
- func (x *DurableTaskEventLogEntryCompletedResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEventLogEntryCompletedResponse) Reset()
- func (x *DurableTaskEventLogEntryCompletedResponse) String() string
- type DurableTaskEventMemoAckResponse
- func (*DurableTaskEventMemoAckResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEventMemoAckResponse) GetMemoAlreadyExisted() bool
- func (x *DurableTaskEventMemoAckResponse) GetMemoResultPayload() []byte
- func (x *DurableTaskEventMemoAckResponse) GetRef() *DurableEventLogEntryRef
- func (*DurableTaskEventMemoAckResponse) ProtoMessage()
- func (x *DurableTaskEventMemoAckResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEventMemoAckResponse) Reset()
- func (x *DurableTaskEventMemoAckResponse) String() string
- type DurableTaskEventTriggerRunsAckResponse
- func (*DurableTaskEventTriggerRunsAckResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEventTriggerRunsAckResponse) GetDurableTaskExternalId() string
- func (x *DurableTaskEventTriggerRunsAckResponse) GetInvocationCount() int32
- func (x *DurableTaskEventTriggerRunsAckResponse) GetRunEntries() []*DurableTaskRunAckEntry
- func (*DurableTaskEventTriggerRunsAckResponse) ProtoMessage()
- func (x *DurableTaskEventTriggerRunsAckResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEventTriggerRunsAckResponse) Reset()
- func (x *DurableTaskEventTriggerRunsAckResponse) String() string
- type DurableTaskEventWaitForAckResponse
- func (*DurableTaskEventWaitForAckResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEventWaitForAckResponse) GetRef() *DurableEventLogEntryRef
- func (*DurableTaskEventWaitForAckResponse) ProtoMessage()
- func (x *DurableTaskEventWaitForAckResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEventWaitForAckResponse) Reset()
- func (x *DurableTaskEventWaitForAckResponse) String() string
- type DurableTaskEvictInvocationRequest
- func (*DurableTaskEvictInvocationRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEvictInvocationRequest) GetDurableTaskExternalId() string
- func (x *DurableTaskEvictInvocationRequest) GetInvocationCount() int32
- func (x *DurableTaskEvictInvocationRequest) GetReason() string
- func (*DurableTaskEvictInvocationRequest) ProtoMessage()
- func (x *DurableTaskEvictInvocationRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEvictInvocationRequest) Reset()
- func (x *DurableTaskEvictInvocationRequest) String() string
- type DurableTaskEvictionAckResponse
- func (*DurableTaskEvictionAckResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskEvictionAckResponse) GetDurableTaskExternalId() string
- func (x *DurableTaskEvictionAckResponse) GetInvocationCount() int32
- func (*DurableTaskEvictionAckResponse) ProtoMessage()
- func (x *DurableTaskEvictionAckResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskEvictionAckResponse) Reset()
- func (x *DurableTaskEvictionAckResponse) String() string
- type DurableTaskMemoRequest
- func (*DurableTaskMemoRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskMemoRequest) GetDurableTaskExternalId() string
- func (x *DurableTaskMemoRequest) GetInvocationCount() int32
- func (x *DurableTaskMemoRequest) GetKey() []byte
- func (x *DurableTaskMemoRequest) GetPayload() []byte
- func (*DurableTaskMemoRequest) ProtoMessage()
- func (x *DurableTaskMemoRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskMemoRequest) Reset()
- func (x *DurableTaskMemoRequest) String() string
- type DurableTaskRequest
- func (*DurableTaskRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskRequest) GetCompleteMemo() *DurableTaskCompleteMemoRequest
- func (x *DurableTaskRequest) GetEvictInvocation() *DurableTaskEvictInvocationRequest
- func (x *DurableTaskRequest) GetMemo() *DurableTaskMemoRequest
- func (m *DurableTaskRequest) GetMessage() isDurableTaskRequest_Message
- func (x *DurableTaskRequest) GetRegisterWorker() *DurableTaskRequestRegisterWorker
- func (x *DurableTaskRequest) GetTriggerRuns() *DurableTaskTriggerRunsRequest
- func (x *DurableTaskRequest) GetWaitFor() *DurableTaskWaitForRequest
- func (x *DurableTaskRequest) GetWorkerStatus() *DurableTaskWorkerStatusRequest
- func (*DurableTaskRequest) ProtoMessage()
- func (x *DurableTaskRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskRequest) Reset()
- func (x *DurableTaskRequest) String() string
- type DurableTaskRequestRegisterWorker
- func (*DurableTaskRequestRegisterWorker) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskRequestRegisterWorker) GetWorkerId() string
- func (*DurableTaskRequestRegisterWorker) ProtoMessage()
- func (x *DurableTaskRequestRegisterWorker) ProtoReflect() protoreflect.Message
- func (x *DurableTaskRequestRegisterWorker) Reset()
- func (x *DurableTaskRequestRegisterWorker) String() string
- type DurableTaskRequest_CompleteMemo
- type DurableTaskRequest_EvictInvocation
- type DurableTaskRequest_Memo
- type DurableTaskRequest_RegisterWorker
- type DurableTaskRequest_TriggerRuns
- type DurableTaskRequest_WaitFor
- type DurableTaskRequest_WorkerStatus
- type DurableTaskResponse
- func (*DurableTaskResponse) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskResponse) GetEntryCompleted() *DurableTaskEventLogEntryCompletedResponse
- func (x *DurableTaskResponse) GetError() *DurableTaskErrorResponse
- func (x *DurableTaskResponse) GetEvictionAck() *DurableTaskEvictionAckResponse
- func (x *DurableTaskResponse) GetMemoAck() *DurableTaskEventMemoAckResponse
- func (m *DurableTaskResponse) GetMessage() isDurableTaskResponse_Message
- func (x *DurableTaskResponse) GetRegisterWorker() *DurableTaskResponseRegisterWorker
- func (x *DurableTaskResponse) GetServerEvict() *DurableTaskServerEvictNotice
- func (x *DurableTaskResponse) GetTriggerRunsAck() *DurableTaskEventTriggerRunsAckResponse
- func (x *DurableTaskResponse) GetWaitForAck() *DurableTaskEventWaitForAckResponse
- func (*DurableTaskResponse) ProtoMessage()
- func (x *DurableTaskResponse) ProtoReflect() protoreflect.Message
- func (x *DurableTaskResponse) Reset()
- func (x *DurableTaskResponse) String() string
- type DurableTaskResponseRegisterWorker
- func (*DurableTaskResponseRegisterWorker) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskResponseRegisterWorker) GetWorkerId() string
- func (*DurableTaskResponseRegisterWorker) ProtoMessage()
- func (x *DurableTaskResponseRegisterWorker) ProtoReflect() protoreflect.Message
- func (x *DurableTaskResponseRegisterWorker) Reset()
- func (x *DurableTaskResponseRegisterWorker) String() string
- type DurableTaskResponse_EntryCompleted
- type DurableTaskResponse_Error
- type DurableTaskResponse_EvictionAck
- type DurableTaskResponse_MemoAck
- type DurableTaskResponse_RegisterWorker
- type DurableTaskResponse_ServerEvict
- type DurableTaskResponse_TriggerRunsAck
- type DurableTaskResponse_WaitForAck
- type DurableTaskRunAckEntry
- func (*DurableTaskRunAckEntry) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskRunAckEntry) GetBranchId() int64
- func (x *DurableTaskRunAckEntry) GetNodeId() int64
- func (x *DurableTaskRunAckEntry) GetWorkflowRunExternalId() string
- func (*DurableTaskRunAckEntry) ProtoMessage()
- func (x *DurableTaskRunAckEntry) ProtoReflect() protoreflect.Message
- func (x *DurableTaskRunAckEntry) Reset()
- func (x *DurableTaskRunAckEntry) String() string
- type DurableTaskServerEvictNotice
- func (*DurableTaskServerEvictNotice) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskServerEvictNotice) GetDurableTaskExternalId() string
- func (x *DurableTaskServerEvictNotice) GetInvocationCount() int32
- func (x *DurableTaskServerEvictNotice) GetReason() string
- func (*DurableTaskServerEvictNotice) ProtoMessage()
- func (x *DurableTaskServerEvictNotice) ProtoReflect() protoreflect.Message
- func (x *DurableTaskServerEvictNotice) Reset()
- func (x *DurableTaskServerEvictNotice) String() string
- type DurableTaskTriggerRunsRequest
- func (*DurableTaskTriggerRunsRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskTriggerRunsRequest) GetDurableTaskExternalId() string
- func (x *DurableTaskTriggerRunsRequest) GetInvocationCount() int32
- func (x *DurableTaskTriggerRunsRequest) GetTriggerOpts() []*TriggerWorkflowRequest
- func (*DurableTaskTriggerRunsRequest) ProtoMessage()
- func (x *DurableTaskTriggerRunsRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskTriggerRunsRequest) Reset()
- func (x *DurableTaskTriggerRunsRequest) String() string
- type DurableTaskWaitForRequest
- func (*DurableTaskWaitForRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskWaitForRequest) GetDurableTaskExternalId() string
- func (x *DurableTaskWaitForRequest) GetInvocationCount() int32
- func (x *DurableTaskWaitForRequest) GetLabel() string
- func (x *DurableTaskWaitForRequest) GetWaitForConditions() *DurableEventListenerConditions
- func (*DurableTaskWaitForRequest) ProtoMessage()
- func (x *DurableTaskWaitForRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskWaitForRequest) Reset()
- func (x *DurableTaskWaitForRequest) String() string
- type DurableTaskWorkerStatusRequest
- func (*DurableTaskWorkerStatusRequest) Descriptor() ([]byte, []int)deprecated
- func (x *DurableTaskWorkerStatusRequest) GetWaitingEntries() []*DurableTaskAwaitedCompletedEntry
- func (x *DurableTaskWorkerStatusRequest) GetWorkerId() string
- func (*DurableTaskWorkerStatusRequest) ProtoMessage()
- func (x *DurableTaskWorkerStatusRequest) ProtoReflect() protoreflect.Message
- func (x *DurableTaskWorkerStatusRequest) Reset()
- func (x *DurableTaskWorkerStatusRequest) String() string
- type GetRunDetailsRequest
- func (*GetRunDetailsRequest) Descriptor() ([]byte, []int)deprecated
- func (x *GetRunDetailsRequest) GetExternalId() string
- func (*GetRunDetailsRequest) ProtoMessage()
- func (x *GetRunDetailsRequest) ProtoReflect() protoreflect.Message
- func (x *GetRunDetailsRequest) Reset()
- func (x *GetRunDetailsRequest) String() string
- type GetRunDetailsResponse
- func (*GetRunDetailsResponse) Descriptor() ([]byte, []int)deprecated
- func (x *GetRunDetailsResponse) GetAdditionalMetadata() []byte
- func (x *GetRunDetailsResponse) GetDone() bool
- func (x *GetRunDetailsResponse) GetInput() []byte
- func (x *GetRunDetailsResponse) GetIsEvicted() bool
- func (x *GetRunDetailsResponse) GetStatus() RunStatus
- func (x *GetRunDetailsResponse) GetTaskRuns() map[string]*TaskRunDetail
- func (*GetRunDetailsResponse) ProtoMessage()
- func (x *GetRunDetailsResponse) ProtoReflect() protoreflect.Message
- func (x *GetRunDetailsResponse) Reset()
- func (x *GetRunDetailsResponse) String() string
- type GetStreamTopicMetadataRequest
- func (*GetStreamTopicMetadataRequest) Descriptor() ([]byte, []int)deprecated
- func (x *GetStreamTopicMetadataRequest) GetNamespace() string
- func (x *GetStreamTopicMetadataRequest) GetTopic() string
- func (*GetStreamTopicMetadataRequest) ProtoMessage()
- func (x *GetStreamTopicMetadataRequest) ProtoReflect() protoreflect.Message
- func (x *GetStreamTopicMetadataRequest) Reset()
- func (x *GetStreamTopicMetadataRequest) String() string
- type IdempotencyCollisionError
- func (*IdempotencyCollisionError) Descriptor() ([]byte, []int)deprecated
- func (x *IdempotencyCollisionError) GetCollidingRunExternalId() string
- func (x *IdempotencyCollisionError) GetExistingRunExternalId() string
- func (*IdempotencyCollisionError) ProtoMessage()
- func (x *IdempotencyCollisionError) ProtoReflect() protoreflect.Message
- func (x *IdempotencyCollisionError) Reset()
- func (x *IdempotencyCollisionError) String() string
- type IdempotencyConfig
- func (*IdempotencyConfig) Descriptor() ([]byte, []int)deprecated
- func (x *IdempotencyConfig) GetExpression() string
- func (x *IdempotencyConfig) GetMethod() IdempotencyMethod
- func (x *IdempotencyConfig) GetTtlMs() int64
- func (*IdempotencyConfig) ProtoMessage()
- func (x *IdempotencyConfig) ProtoReflect() protoreflect.Message
- func (x *IdempotencyConfig) Reset()
- func (x *IdempotencyConfig) String() string
- type IdempotencyMethod
- func (IdempotencyMethod) Descriptor() protoreflect.EnumDescriptor
- func (x IdempotencyMethod) Enum() *IdempotencyMethod
- func (IdempotencyMethod) EnumDescriptor() ([]byte, []int)deprecated
- func (x IdempotencyMethod) Number() protoreflect.EnumNumber
- func (x IdempotencyMethod) String() string
- func (IdempotencyMethod) Type() protoreflect.EnumType
- type ListenForDurableEventRequest
- func (*ListenForDurableEventRequest) Descriptor() ([]byte, []int)deprecated
- func (x *ListenForDurableEventRequest) GetSignalKey() string
- func (x *ListenForDurableEventRequest) GetTaskId() string
- func (*ListenForDurableEventRequest) ProtoMessage()
- func (x *ListenForDurableEventRequest) ProtoReflect() protoreflect.Message
- func (x *ListenForDurableEventRequest) Reset()
- func (x *ListenForDurableEventRequest) String() string
- type OperatorActionsAck
- func (*OperatorActionsAck) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorActionsAck) GetSequence() uint64
- func (*OperatorActionsAck) ProtoMessage()
- func (x *OperatorActionsAck) ProtoReflect() protoreflect.Message
- func (x *OperatorActionsAck) Reset()
- func (x *OperatorActionsAck) String() string
- type OperatorActionsDelta
- func (*OperatorActionsDelta) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorActionsDelta) GetAdd() []string
- func (x *OperatorActionsDelta) GetRemove() []string
- func (x *OperatorActionsDelta) GetSequence() uint64
- func (*OperatorActionsDelta) ProtoMessage()
- func (x *OperatorActionsDelta) ProtoReflect() protoreflect.Message
- func (x *OperatorActionsDelta) Reset()
- func (x *OperatorActionsDelta) String() string
- type OperatorHeartbeat
- func (*OperatorHeartbeat) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorHeartbeat) GetHeartbeatAt() *timestamppb.Timestamp
- func (*OperatorHeartbeat) ProtoMessage()
- func (x *OperatorHeartbeat) ProtoReflect() protoreflect.Message
- func (x *OperatorHeartbeat) Reset()
- func (x *OperatorHeartbeat) String() string
- type OperatorListenRequest
- func (*OperatorListenRequest) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorListenRequest) GetActions() *OperatorActionsDelta
- func (x *OperatorListenRequest) GetHeartbeat() *OperatorHeartbeat
- func (m *OperatorListenRequest) GetMessage() isOperatorListenRequest_Message
- func (x *OperatorListenRequest) GetPause() *OperatorPause
- func (x *OperatorListenRequest) GetStart() *OperatorListenStart
- func (*OperatorListenRequest) ProtoMessage()
- func (x *OperatorListenRequest) ProtoReflect() protoreflect.Message
- func (x *OperatorListenRequest) Reset()
- func (x *OperatorListenRequest) String() string
- type OperatorListenRequest_Actions
- type OperatorListenRequest_Heartbeat
- type OperatorListenRequest_Pause
- type OperatorListenRequest_Start
- type OperatorListenResponse
- func (*OperatorListenResponse) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorListenResponse) GetAck() *OperatorActionsAck
- func (x *OperatorListenResponse) GetAction() *contracts.AssignedAction
- func (m *OperatorListenResponse) GetMessage() isOperatorListenResponse_Message
- func (x *OperatorListenResponse) GetPauseAck() *OperatorPauseAck
- func (*OperatorListenResponse) ProtoMessage()
- func (x *OperatorListenResponse) ProtoReflect() protoreflect.Message
- func (x *OperatorListenResponse) Reset()
- func (x *OperatorListenResponse) String() string
- type OperatorListenResponse_Ack
- type OperatorListenResponse_Action
- type OperatorListenResponse_PauseAck
- type OperatorListenStart
- func (*OperatorListenStart) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorListenStart) GetWorkerId() string
- func (*OperatorListenStart) ProtoMessage()
- func (x *OperatorListenStart) ProtoReflect() protoreflect.Message
- func (x *OperatorListenStart) Reset()
- func (x *OperatorListenStart) String() string
- type OperatorPause
- type OperatorPauseAck
- type OperatorRegisterRequest
- func (*OperatorRegisterRequest) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorRegisterRequest) GetLabels() map[string]*contracts.WorkerLabels
- func (x *OperatorRegisterRequest) GetName() string
- func (x *OperatorRegisterRequest) GetRuntimeInfo() *contracts.RuntimeInfo
- func (x *OperatorRegisterRequest) GetSlotConfig() map[string]int32
- func (x *OperatorRegisterRequest) GetWorkerId() string
- func (*OperatorRegisterRequest) ProtoMessage()
- func (x *OperatorRegisterRequest) ProtoReflect() protoreflect.Message
- func (x *OperatorRegisterRequest) Reset()
- func (x *OperatorRegisterRequest) String() string
- type OperatorRegisterResponse
- func (*OperatorRegisterResponse) Descriptor() ([]byte, []int)deprecated
- func (x *OperatorRegisterResponse) GetOperatorId() string
- func (x *OperatorRegisterResponse) GetResumed() bool
- func (x *OperatorRegisterResponse) GetTenantId() string
- func (x *OperatorRegisterResponse) GetWorkerId() string
- func (*OperatorRegisterResponse) ProtoMessage()
- func (x *OperatorRegisterResponse) ProtoReflect() protoreflect.Message
- func (x *OperatorRegisterResponse) Reset()
- func (x *OperatorRegisterResponse) String() string
- type OperatorServiceClient
- type OperatorServiceServer
- type OperatorService_DurableTaskClient
- type OperatorService_DurableTaskServer
- type OperatorService_ListenClient
- type OperatorService_ListenServer
- type ParentOverrideMatchCondition
- func (*ParentOverrideMatchCondition) Descriptor() ([]byte, []int)deprecated
- func (x *ParentOverrideMatchCondition) GetBase() *BaseMatchCondition
- func (x *ParentOverrideMatchCondition) GetParentReadableId() string
- func (*ParentOverrideMatchCondition) ProtoMessage()
- func (x *ParentOverrideMatchCondition) ProtoReflect() protoreflect.Message
- func (x *ParentOverrideMatchCondition) Reset()
- func (x *ParentOverrideMatchCondition) String() string
- type PublishStreamMessageRequest
- func (*PublishStreamMessageRequest) Descriptor() ([]byte, []int)deprecated
- func (x *PublishStreamMessageRequest) GetNamespace() string
- func (x *PublishStreamMessageRequest) GetPayload() []byte
- func (x *PublishStreamMessageRequest) GetPayloadRef() string
- func (x *PublishStreamMessageRequest) GetProducerId() string
- func (x *PublishStreamMessageRequest) GetProducerSeq() int64
- func (x *PublishStreamMessageRequest) GetTopic() string
- func (*PublishStreamMessageRequest) ProtoMessage()
- func (x *PublishStreamMessageRequest) ProtoReflect() protoreflect.Message
- func (x *PublishStreamMessageRequest) Reset()
- func (x *PublishStreamMessageRequest) String() string
- type PublishStreamMessageResponse
- func (*PublishStreamMessageResponse) Descriptor() ([]byte, []int)deprecated
- func (*PublishStreamMessageResponse) ProtoMessage()
- func (x *PublishStreamMessageResponse) ProtoReflect() protoreflect.Message
- func (x *PublishStreamMessageResponse) Reset()
- func (x *PublishStreamMessageResponse) String() string
- type RateLimitDuration
- func (RateLimitDuration) Descriptor() protoreflect.EnumDescriptor
- func (x RateLimitDuration) Enum() *RateLimitDuration
- func (RateLimitDuration) EnumDescriptor() ([]byte, []int)deprecated
- func (x RateLimitDuration) Number() protoreflect.EnumNumber
- func (x RateLimitDuration) String() string
- func (RateLimitDuration) Type() protoreflect.EnumType
- type RegisterDurableEventRequest
- func (*RegisterDurableEventRequest) Descriptor() ([]byte, []int)deprecated
- func (x *RegisterDurableEventRequest) GetConditions() *DurableEventListenerConditions
- func (x *RegisterDurableEventRequest) GetSignalKey() string
- func (x *RegisterDurableEventRequest) GetTaskId() string
- func (*RegisterDurableEventRequest) ProtoMessage()
- func (x *RegisterDurableEventRequest) ProtoReflect() protoreflect.Message
- func (x *RegisterDurableEventRequest) Reset()
- func (x *RegisterDurableEventRequest) String() string
- type RegisterDurableEventResponse
- func (*RegisterDurableEventResponse) Descriptor() ([]byte, []int)deprecated
- func (*RegisterDurableEventResponse) ProtoMessage()
- func (x *RegisterDurableEventResponse) ProtoReflect() protoreflect.Message
- func (x *RegisterDurableEventResponse) Reset()
- func (x *RegisterDurableEventResponse) String() string
- type ReplayTasksRequest
- func (*ReplayTasksRequest) Descriptor() ([]byte, []int)deprecated
- func (x *ReplayTasksRequest) GetExternalIds() []string
- func (x *ReplayTasksRequest) GetFilter() *TasksFilter
- func (*ReplayTasksRequest) ProtoMessage()
- func (x *ReplayTasksRequest) ProtoReflect() protoreflect.Message
- func (x *ReplayTasksRequest) Reset()
- func (x *ReplayTasksRequest) String() string
- type ReplayTasksResponse
- func (*ReplayTasksResponse) Descriptor() ([]byte, []int)deprecated
- func (x *ReplayTasksResponse) GetReplayedTasks() []string
- func (*ReplayTasksResponse) ProtoMessage()
- func (x *ReplayTasksResponse) ProtoReflect() protoreflect.Message
- func (x *ReplayTasksResponse) Reset()
- func (x *ReplayTasksResponse) String() string
- type RunStatus
- type SleepMatchCondition
- func (*SleepMatchCondition) Descriptor() ([]byte, []int)deprecated
- func (x *SleepMatchCondition) GetBase() *BaseMatchCondition
- func (x *SleepMatchCondition) GetSleepFor() string
- func (*SleepMatchCondition) ProtoMessage()
- func (x *SleepMatchCondition) ProtoReflect() protoreflect.Message
- func (x *SleepMatchCondition) Reset()
- func (x *SleepMatchCondition) String() string
- type StickyStrategy
- func (StickyStrategy) Descriptor() protoreflect.EnumDescriptor
- func (x StickyStrategy) Enum() *StickyStrategy
- func (StickyStrategy) EnumDescriptor() ([]byte, []int)deprecated
- func (x StickyStrategy) Number() protoreflect.EnumNumber
- func (x StickyStrategy) String() string
- func (StickyStrategy) Type() protoreflect.EnumType
- type StreamEntry
- func (*StreamEntry) Descriptor() ([]byte, []int)deprecated
- func (x *StreamEntry) GetCreatedAt() *timestamppb.Timestamp
- func (x *StreamEntry) GetCursor() string
- func (x *StreamEntry) GetPayload() []byte
- func (x *StreamEntry) GetPayloadRef() string
- func (*StreamEntry) ProtoMessage()
- func (x *StreamEntry) ProtoReflect() protoreflect.Message
- func (x *StreamEntry) Reset()
- func (x *StreamEntry) String() string
- type StreamMessage
- func (*StreamMessage) Descriptor() ([]byte, []int)deprecated
- func (x *StreamMessage) GetCursor() string
- func (x *StreamMessage) GetEntries() []*StreamEntry
- func (x *StreamMessage) GetHangup() bool
- func (*StreamMessage) ProtoMessage()
- func (x *StreamMessage) ProtoReflect() protoreflect.Message
- func (x *StreamMessage) Reset()
- func (x *StreamMessage) String() string
- type StreamTopicMetadata
- func (*StreamTopicMetadata) Descriptor() ([]byte, []int)deprecated
- func (x *StreamTopicMetadata) GetLastPublishedAt() *timestamppb.Timestamp
- func (x *StreamTopicMetadata) GetLatestCursor() string
- func (x *StreamTopicMetadata) GetMessageCount() int64
- func (x *StreamTopicMetadata) GetNamespace() string
- func (x *StreamTopicMetadata) GetTenantId() string
- func (x *StreamTopicMetadata) GetTopic() string
- func (*StreamTopicMetadata) ProtoMessage()
- func (x *StreamTopicMetadata) ProtoReflect() protoreflect.Message
- func (x *StreamTopicMetadata) Reset()
- func (x *StreamTopicMetadata) String() string
- type SubscribeStreamRequest
- func (*SubscribeStreamRequest) Descriptor() ([]byte, []int)deprecated
- func (x *SubscribeStreamRequest) GetCursor() string
- func (x *SubscribeStreamRequest) GetNamespace() string
- func (x *SubscribeStreamRequest) GetTopic() string
- func (*SubscribeStreamRequest) ProtoMessage()
- func (x *SubscribeStreamRequest) ProtoReflect() protoreflect.Message
- func (x *SubscribeStreamRequest) Reset()
- func (x *SubscribeStreamRequest) String() string
- type TaskBatchConfig
- func (*TaskBatchConfig) Descriptor() ([]byte, []int)deprecated
- func (x *TaskBatchConfig) GetBatchGroupKey() string
- func (x *TaskBatchConfig) GetBatchGroupMaxRuns() int32
- func (x *TaskBatchConfig) GetBatchMaxIntervalMs() int32
- func (x *TaskBatchConfig) GetBatchMaxSize() int32
- func (x *TaskBatchConfig) GetBroadcastOutput() bool
- func (*TaskBatchConfig) ProtoMessage()
- func (x *TaskBatchConfig) ProtoReflect() protoreflect.Message
- func (x *TaskBatchConfig) Reset()
- func (x *TaskBatchConfig) String() string
- type TaskConditions
- func (*TaskConditions) Descriptor() ([]byte, []int)deprecated
- func (x *TaskConditions) GetParentOverrideConditions() []*ParentOverrideMatchCondition
- func (x *TaskConditions) GetSleepConditions() []*SleepMatchCondition
- func (x *TaskConditions) GetUserEventConditions() []*UserEventMatchCondition
- func (*TaskConditions) ProtoMessage()
- func (x *TaskConditions) ProtoReflect() protoreflect.Message
- func (x *TaskConditions) Reset()
- func (x *TaskConditions) String() string
- type TaskRunDetail
- func (*TaskRunDetail) Descriptor() ([]byte, []int)deprecated
- func (x *TaskRunDetail) GetError() string
- func (x *TaskRunDetail) GetExternalId() string
- func (x *TaskRunDetail) GetIsEvicted() bool
- func (x *TaskRunDetail) GetOutput() []byte
- func (x *TaskRunDetail) GetReadableId() string
- func (x *TaskRunDetail) GetStatus() RunStatus
- func (*TaskRunDetail) ProtoMessage()
- func (x *TaskRunDetail) ProtoReflect() protoreflect.Message
- func (x *TaskRunDetail) Reset()
- func (x *TaskRunDetail) String() string
- type TasksFilter
- func (*TasksFilter) Descriptor() ([]byte, []int)deprecated
- func (x *TasksFilter) GetAdditionalMetadata() []string
- func (x *TasksFilter) GetSince() *timestamppb.Timestamp
- func (x *TasksFilter) GetStatuses() []string
- func (x *TasksFilter) GetUntil() *timestamppb.Timestamp
- func (x *TasksFilter) GetWorkflowIds() []string
- func (*TasksFilter) ProtoMessage()
- func (x *TasksFilter) ProtoReflect() protoreflect.Message
- func (x *TasksFilter) Reset()
- func (x *TasksFilter) String() string
- type TriggerWorkflowRequest
- func (*TriggerWorkflowRequest) Descriptor() ([]byte, []int)deprecated
- func (x *TriggerWorkflowRequest) GetAdditionalMetadata() string
- func (x *TriggerWorkflowRequest) GetChildIndex() int32
- func (x *TriggerWorkflowRequest) GetChildKey() string
- func (x *TriggerWorkflowRequest) GetDesiredWorkerId() string
- func (x *TriggerWorkflowRequest) GetDesiredWorkerLabels() map[string]*DesiredWorkerLabels
- func (x *TriggerWorkflowRequest) GetInput() string
- func (x *TriggerWorkflowRequest) GetName() string
- func (x *TriggerWorkflowRequest) GetParentId() string
- func (x *TriggerWorkflowRequest) GetParentTaskRunExternalId() string
- func (x *TriggerWorkflowRequest) GetPriority() int32
- func (*TriggerWorkflowRequest) ProtoMessage()
- func (x *TriggerWorkflowRequest) ProtoReflect() protoreflect.Message
- func (x *TriggerWorkflowRequest) Reset()
- func (x *TriggerWorkflowRequest) String() string
- type TriggerWorkflowRunRequest
- func (*TriggerWorkflowRunRequest) Descriptor() ([]byte, []int)deprecated
- func (x *TriggerWorkflowRunRequest) GetAdditionalMetadata() []byte
- func (x *TriggerWorkflowRunRequest) GetDesiredWorkerLabels() map[string]*DesiredWorkerLabels
- func (x *TriggerWorkflowRunRequest) GetInput() []byte
- func (x *TriggerWorkflowRunRequest) GetPriority() int32
- func (x *TriggerWorkflowRunRequest) GetWorkflowName() string
- func (*TriggerWorkflowRunRequest) ProtoMessage()
- func (x *TriggerWorkflowRunRequest) ProtoReflect() protoreflect.Message
- func (x *TriggerWorkflowRunRequest) Reset()
- func (x *TriggerWorkflowRunRequest) String() string
- type TriggerWorkflowRunResponse
- func (*TriggerWorkflowRunResponse) Descriptor() ([]byte, []int)deprecated
- func (x *TriggerWorkflowRunResponse) GetExternalId() string
- func (*TriggerWorkflowRunResponse) ProtoMessage()
- func (x *TriggerWorkflowRunResponse) ProtoReflect() protoreflect.Message
- func (x *TriggerWorkflowRunResponse) Reset()
- func (x *TriggerWorkflowRunResponse) String() string
- type UnimplementedAdminServiceServer
- func (UnimplementedAdminServiceServer) BranchDurableTask(context.Context, *BranchDurableTaskRequest) (*BranchDurableTaskResponse, error)
- func (UnimplementedAdminServiceServer) CancelTasks(context.Context, *CancelTasksRequest) (*CancelTasksResponse, error)
- func (UnimplementedAdminServiceServer) GetRunDetails(context.Context, *GetRunDetailsRequest) (*GetRunDetailsResponse, error)
- func (UnimplementedAdminServiceServer) PutWorkflow(context.Context, *CreateWorkflowVersionRequest) (*CreateWorkflowVersionResponse, error)
- func (UnimplementedAdminServiceServer) ReplayTasks(context.Context, *ReplayTasksRequest) (*ReplayTasksResponse, error)
- func (UnimplementedAdminServiceServer) TriggerWorkflowRun(context.Context, *TriggerWorkflowRunRequest) (*TriggerWorkflowRunResponse, error)
- type UnimplementedOperatorServiceServer
- func (UnimplementedOperatorServiceServer) DurableTask(OperatorService_DurableTaskServer) error
- func (UnimplementedOperatorServiceServer) Listen(OperatorService_ListenServer) error
- func (UnimplementedOperatorServiceServer) Register(context.Context, *OperatorRegisterRequest) (*OperatorRegisterResponse, error)
- func (UnimplementedOperatorServiceServer) SendStepActionEvent(context.Context, *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)
- type UnimplementedV1DispatcherServer
- func (UnimplementedV1DispatcherServer) DurableTask(V1Dispatcher_DurableTaskServer) error
- func (UnimplementedV1DispatcherServer) ListenForDurableEvent(V1Dispatcher_ListenForDurableEventServer) error
- func (UnimplementedV1DispatcherServer) RegisterDurableEvent(context.Context, *RegisterDurableEventRequest) (*RegisterDurableEventResponse, error)
- type UnimplementedV1StreamsServer
- func (UnimplementedV1StreamsServer) GetTopicMetadata(context.Context, *GetStreamTopicMetadataRequest) (*StreamTopicMetadata, error)
- func (UnimplementedV1StreamsServer) Publish(context.Context, *PublishStreamMessageRequest) (*PublishStreamMessageResponse, error)
- func (UnimplementedV1StreamsServer) Subscribe(*SubscribeStreamRequest, V1Streams_SubscribeServer) error
- type UnsafeAdminServiceServer
- type UnsafeOperatorServiceServer
- type UnsafeV1DispatcherServer
- type UnsafeV1StreamsServer
- type UserEventMatchCondition
- func (*UserEventMatchCondition) Descriptor() ([]byte, []int)deprecated
- func (x *UserEventMatchCondition) GetBase() *BaseMatchCondition
- func (x *UserEventMatchCondition) GetConsiderEventsSince() *timestamppb.Timestamp
- func (x *UserEventMatchCondition) GetEventScope() string
- func (x *UserEventMatchCondition) GetUserEventKey() string
- func (*UserEventMatchCondition) ProtoMessage()
- func (x *UserEventMatchCondition) ProtoReflect() protoreflect.Message
- func (x *UserEventMatchCondition) Reset()
- func (x *UserEventMatchCondition) String() string
- type V1DispatcherClient
- type V1DispatcherServer
- type V1Dispatcher_DurableTaskClient
- type V1Dispatcher_DurableTaskServer
- type V1Dispatcher_ListenForDurableEventClient
- type V1Dispatcher_ListenForDurableEventServer
- type V1StreamsClient
- type V1StreamsServer
- type V1Streams_SubscribeClient
- type V1Streams_SubscribeServer
- type WorkerLabelComparator
- func (WorkerLabelComparator) Descriptor() protoreflect.EnumDescriptor
- func (x WorkerLabelComparator) Enum() *WorkerLabelComparator
- func (WorkerLabelComparator) EnumDescriptor() ([]byte, []int)deprecated
- func (x WorkerLabelComparator) Number() protoreflect.EnumNumber
- func (x WorkerLabelComparator) String() string
- func (WorkerLabelComparator) Type() protoreflect.EnumType
Constants ¶
This section is empty.
Variables ¶
var ( Action_name = map[int32]string{ 0: "CREATE", 1: "QUEUE", 2: "CANCEL", 3: "SKIP", } Action_value = map[string]int32{ "CREATE": 0, "QUEUE": 1, "CANCEL": 2, "SKIP": 3, } )
Enum value maps for Action.
var ( DurableTaskErrorType_name = map[int32]string{ 0: "DURABLE_TASK_ERROR_TYPE_UNSPECIFIED", 1: "DURABLE_TASK_ERROR_TYPE_NONDETERMINISM", } DurableTaskErrorType_value = map[string]int32{ "DURABLE_TASK_ERROR_TYPE_UNSPECIFIED": 0, "DURABLE_TASK_ERROR_TYPE_NONDETERMINISM": 1, } )
Enum value maps for DurableTaskErrorType.
var ( WorkerLabelComparator_name = map[int32]string{ 0: "EQUAL", 1: "NOT_EQUAL", 2: "GREATER_THAN", 3: "GREATER_THAN_OR_EQUAL", 4: "LESS_THAN", 5: "LESS_THAN_OR_EQUAL", } WorkerLabelComparator_value = map[string]int32{ "EQUAL": 0, "NOT_EQUAL": 1, "GREATER_THAN": 2, "GREATER_THAN_OR_EQUAL": 3, "LESS_THAN": 4, "LESS_THAN_OR_EQUAL": 5, } )
Enum value maps for WorkerLabelComparator.
var ( StickyStrategy_name = map[int32]string{ 0: "SOFT", 1: "HARD", } StickyStrategy_value = map[string]int32{ "SOFT": 0, "HARD": 1, } )
Enum value maps for StickyStrategy.
var ( RateLimitDuration_name = map[int32]string{ 0: "SECOND", 1: "MINUTE", 2: "HOUR", 3: "DAY", 4: "WEEK", 5: "MONTH", 6: "YEAR", } RateLimitDuration_value = map[string]int32{ "SECOND": 0, "MINUTE": 1, "HOUR": 2, "DAY": 3, "WEEK": 4, "MONTH": 5, "YEAR": 6, } )
Enum value maps for RateLimitDuration.
var ( RunStatus_name = map[int32]string{ 0: "QUEUED", 1: "RUNNING", 2: "COMPLETED", 3: "FAILED", 4: "CANCELLED", 5: "EVICTED", } RunStatus_value = map[string]int32{ "QUEUED": 0, "RUNNING": 1, "COMPLETED": 2, "FAILED": 3, "CANCELLED": 4, "EVICTED": 5, } )
Enum value maps for RunStatus.
var ( IdempotencyMethod_name = map[int32]string{ 0: "TTL", 1: "STATUS", } IdempotencyMethod_value = map[string]int32{ "TTL": 0, "STATUS": 1, } )
Enum value maps for IdempotencyMethod.
var ( ConcurrencyLimitStrategy_name = map[int32]string{ 0: "CANCEL_IN_PROGRESS", 1: "DROP_NEWEST", 2: "QUEUE_NEWEST", 3: "GROUP_ROUND_ROBIN", 4: "CANCEL_NEWEST", 5: "CANCEL_QUEUED_EXCEPT_NEWEST", 6: "CANCEL_QUEUED_EXCEPT_OLDEST", } ConcurrencyLimitStrategy_value = map[string]int32{ "CANCEL_IN_PROGRESS": 0, "DROP_NEWEST": 1, "QUEUE_NEWEST": 2, "GROUP_ROUND_ROBIN": 3, "CANCEL_NEWEST": 4, "CANCEL_QUEUED_EXCEPT_NEWEST": 5, "CANCEL_QUEUED_EXCEPT_OLDEST": 6, } )
Enum value maps for ConcurrencyLimitStrategy.
var AdminService_ServiceDesc = grpc.ServiceDesc{ ServiceName: "v1.AdminService", HandlerType: (*AdminServiceServer)(nil), Methods: []grpc.MethodDesc{ { MethodName: "PutWorkflow", Handler: _AdminService_PutWorkflow_Handler, }, { MethodName: "CancelTasks", Handler: _AdminService_CancelTasks_Handler, }, { MethodName: "ReplayTasks", Handler: _AdminService_ReplayTasks_Handler, }, { MethodName: "TriggerWorkflowRun", Handler: _AdminService_TriggerWorkflowRun_Handler, }, { MethodName: "GetRunDetails", Handler: _AdminService_GetRunDetails_Handler, }, { MethodName: "BranchDurableTask", Handler: _AdminService_BranchDurableTask_Handler, }, }, Streams: []grpc.StreamDesc{}, Metadata: "v1/workflows.proto", }
AdminService_ServiceDesc is the grpc.ServiceDesc for AdminService service. It's only intended for direct use with grpc.RegisterService, and not to be introspected or modified (even as a copy)
var File_v1_dispatcher_proto protoreflect.FileDescriptor
var File_v1_operator_proto protoreflect.FileDescriptor
var File_v1_streams_proto protoreflect.FileDescriptor
var File_v1_workflows_proto protoreflect.FileDescriptor
var OperatorService_ServiceDesc = grpc.ServiceDesc{ ServiceName: "v1.OperatorService", HandlerType: (*OperatorServiceServer)(nil), Methods: []grpc.MethodDesc{ { MethodName: "Register", Handler: _OperatorService_Register_Handler, }, { MethodName: "SendStepActionEvent", Handler: _OperatorService_SendStepActionEvent_Handler, }, }, Streams: []grpc.StreamDesc{ { StreamName: "Listen", Handler: _OperatorService_Listen_Handler, ServerStreams: true, ClientStreams: true, }, { StreamName: "DurableTask", Handler: _OperatorService_DurableTask_Handler, ServerStreams: true, ClientStreams: true, }, }, Metadata: "v1/operator.proto", }
OperatorService_ServiceDesc is the grpc.ServiceDesc for OperatorService service. It's only intended for direct use with grpc.RegisterService, and not to be introspected or modified (even as a copy)
var V1Dispatcher_ServiceDesc = grpc.ServiceDesc{ ServiceName: "v1.V1Dispatcher", HandlerType: (*V1DispatcherServer)(nil), Methods: []grpc.MethodDesc{ { MethodName: "RegisterDurableEvent", Handler: _V1Dispatcher_RegisterDurableEvent_Handler, }, }, Streams: []grpc.StreamDesc{ { StreamName: "DurableTask", Handler: _V1Dispatcher_DurableTask_Handler, ServerStreams: true, ClientStreams: true, }, { StreamName: "ListenForDurableEvent", Handler: _V1Dispatcher_ListenForDurableEvent_Handler, ServerStreams: true, ClientStreams: true, }, }, Metadata: "v1/dispatcher.proto", }
V1Dispatcher_ServiceDesc is the grpc.ServiceDesc for V1Dispatcher service. It's only intended for direct use with grpc.RegisterService, and not to be introspected or modified (even as a copy)
var V1Streams_ServiceDesc = grpc.ServiceDesc{ ServiceName: "v1.V1Streams", HandlerType: (*V1StreamsServer)(nil), Methods: []grpc.MethodDesc{ { MethodName: "Publish", Handler: _V1Streams_Publish_Handler, }, { MethodName: "GetTopicMetadata", Handler: _V1Streams_GetTopicMetadata_Handler, }, }, Streams: []grpc.StreamDesc{ { StreamName: "Subscribe", Handler: _V1Streams_Subscribe_Handler, ServerStreams: true, }, }, Metadata: "v1/streams.proto", }
V1Streams_ServiceDesc is the grpc.ServiceDesc for V1Streams service. It's only intended for direct use with grpc.RegisterService, and not to be introspected or modified (even as a copy)
Functions ¶
func RegisterAdminServiceServer ¶
func RegisterAdminServiceServer(s grpc.ServiceRegistrar, srv AdminServiceServer)
func RegisterOperatorServiceServer ¶ added in v0.109.0
func RegisterOperatorServiceServer(s grpc.ServiceRegistrar, srv OperatorServiceServer)
func RegisterV1DispatcherServer ¶
func RegisterV1DispatcherServer(s grpc.ServiceRegistrar, srv V1DispatcherServer)
func RegisterV1StreamsServer ¶ added in v0.110.0
func RegisterV1StreamsServer(s grpc.ServiceRegistrar, srv V1StreamsServer)
Types ¶
type Action ¶
type Action int32
func (Action) Descriptor ¶
func (Action) Descriptor() protoreflect.EnumDescriptor
func (Action) EnumDescriptor
deprecated
func (Action) Number ¶
func (x Action) Number() protoreflect.EnumNumber
func (Action) Type ¶
func (Action) Type() protoreflect.EnumType
type AdminServiceClient ¶
type AdminServiceClient interface {
PutWorkflow(ctx context.Context, in *CreateWorkflowVersionRequest, opts ...grpc.CallOption) (*CreateWorkflowVersionResponse, error)
CancelTasks(ctx context.Context, in *CancelTasksRequest, opts ...grpc.CallOption) (*CancelTasksResponse, error)
ReplayTasks(ctx context.Context, in *ReplayTasksRequest, opts ...grpc.CallOption) (*ReplayTasksResponse, error)
TriggerWorkflowRun(ctx context.Context, in *TriggerWorkflowRunRequest, opts ...grpc.CallOption) (*TriggerWorkflowRunResponse, error)
GetRunDetails(ctx context.Context, in *GetRunDetailsRequest, opts ...grpc.CallOption) (*GetRunDetailsResponse, error)
BranchDurableTask(ctx context.Context, in *BranchDurableTaskRequest, opts ...grpc.CallOption) (*BranchDurableTaskResponse, error)
}
AdminServiceClient is the client API for AdminService service.
For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
func NewAdminServiceClient ¶
func NewAdminServiceClient(cc grpc.ClientConnInterface) AdminServiceClient
type AdminServiceServer ¶
type AdminServiceServer interface {
PutWorkflow(context.Context, *CreateWorkflowVersionRequest) (*CreateWorkflowVersionResponse, error)
CancelTasks(context.Context, *CancelTasksRequest) (*CancelTasksResponse, error)
ReplayTasks(context.Context, *ReplayTasksRequest) (*ReplayTasksResponse, error)
TriggerWorkflowRun(context.Context, *TriggerWorkflowRunRequest) (*TriggerWorkflowRunResponse, error)
GetRunDetails(context.Context, *GetRunDetailsRequest) (*GetRunDetailsResponse, error)
BranchDurableTask(context.Context, *BranchDurableTaskRequest) (*BranchDurableTaskResponse, error)
// contains filtered or unexported methods
}
AdminServiceServer is the server API for AdminService service. All implementations must embed UnimplementedAdminServiceServer for forward compatibility
type BaseMatchCondition ¶
type BaseMatchCondition struct {
ReadableDataKey string `protobuf:"bytes,1,opt,name=readable_data_key,json=readableDataKey,proto3" json:"readable_data_key,omitempty"`
Action Action `protobuf:"varint,2,opt,name=action,proto3,enum=v1.Action" json:"action,omitempty"`
OrGroupId string `protobuf:"bytes,3,opt,name=or_group_id,json=orGroupId,proto3" json:"or_group_id,omitempty"` // a UUID defining the OR group for this condition
Expression string `protobuf:"bytes,4,opt,name=expression,proto3" json:"expression,omitempty"`
// contains filtered or unexported fields
}
func (*BaseMatchCondition) Descriptor
deprecated
func (*BaseMatchCondition) Descriptor() ([]byte, []int)
Deprecated: Use BaseMatchCondition.ProtoReflect.Descriptor instead.
func (*BaseMatchCondition) GetAction ¶
func (x *BaseMatchCondition) GetAction() Action
func (*BaseMatchCondition) GetExpression ¶
func (x *BaseMatchCondition) GetExpression() string
func (*BaseMatchCondition) GetOrGroupId ¶
func (x *BaseMatchCondition) GetOrGroupId() string
func (*BaseMatchCondition) GetReadableDataKey ¶
func (x *BaseMatchCondition) GetReadableDataKey() string
func (*BaseMatchCondition) ProtoMessage ¶
func (*BaseMatchCondition) ProtoMessage()
func (*BaseMatchCondition) ProtoReflect ¶
func (x *BaseMatchCondition) ProtoReflect() protoreflect.Message
func (*BaseMatchCondition) Reset ¶
func (x *BaseMatchCondition) Reset()
func (*BaseMatchCondition) String ¶
func (x *BaseMatchCondition) String() string
type BranchDurableTaskRequest ¶ added in v0.80.0
type BranchDurableTaskRequest struct {
TaskExternalId string `protobuf:"bytes,1,opt,name=task_external_id,json=taskExternalId,proto3" json:"task_external_id,omitempty"` // (required) the external id (uuid) of the durable task
NodeId int64 `protobuf:"varint,2,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` // (required) the node id to branch from
BranchId int64 `protobuf:"varint,3,opt,name=branch_id,json=branchId,proto3" json:"branch_id,omitempty"` // (required) the branch id to branch from
// contains filtered or unexported fields
}
func (*BranchDurableTaskRequest) Descriptor
deprecated
added in
v0.80.0
func (*BranchDurableTaskRequest) Descriptor() ([]byte, []int)
Deprecated: Use BranchDurableTaskRequest.ProtoReflect.Descriptor instead.
func (*BranchDurableTaskRequest) GetBranchId ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) GetBranchId() int64
func (*BranchDurableTaskRequest) GetNodeId ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) GetNodeId() int64
func (*BranchDurableTaskRequest) GetTaskExternalId ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) GetTaskExternalId() string
func (*BranchDurableTaskRequest) ProtoMessage ¶ added in v0.80.0
func (*BranchDurableTaskRequest) ProtoMessage()
func (*BranchDurableTaskRequest) ProtoReflect ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) ProtoReflect() protoreflect.Message
func (*BranchDurableTaskRequest) Reset ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) Reset()
func (*BranchDurableTaskRequest) String ¶ added in v0.80.0
func (x *BranchDurableTaskRequest) String() string
type BranchDurableTaskResponse ¶ added in v0.80.0
type BranchDurableTaskResponse struct {
TaskExternalId string `protobuf:"bytes,1,opt,name=task_external_id,json=taskExternalId,proto3" json:"task_external_id,omitempty"` // the external id of the durable task
NodeId int64 `protobuf:"varint,2,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` // the node id of the new entry
BranchId int64 `protobuf:"varint,3,opt,name=branch_id,json=branchId,proto3" json:"branch_id,omitempty"` // the branch id of the new entry
// contains filtered or unexported fields
}
func (*BranchDurableTaskResponse) Descriptor
deprecated
added in
v0.80.0
func (*BranchDurableTaskResponse) Descriptor() ([]byte, []int)
Deprecated: Use BranchDurableTaskResponse.ProtoReflect.Descriptor instead.
func (*BranchDurableTaskResponse) GetBranchId ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) GetBranchId() int64
func (*BranchDurableTaskResponse) GetNodeId ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) GetNodeId() int64
func (*BranchDurableTaskResponse) GetTaskExternalId ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) GetTaskExternalId() string
func (*BranchDurableTaskResponse) ProtoMessage ¶ added in v0.80.0
func (*BranchDurableTaskResponse) ProtoMessage()
func (*BranchDurableTaskResponse) ProtoReflect ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) ProtoReflect() protoreflect.Message
func (*BranchDurableTaskResponse) Reset ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) Reset()
func (*BranchDurableTaskResponse) String ¶ added in v0.80.0
func (x *BranchDurableTaskResponse) String() string
type BulkTriggerIdempotencyCollisionError ¶ added in v0.95.0
type BulkTriggerIdempotencyCollisionError struct {
SuccessfulWorkflowRunExternalIds []string `` // the external IDs of the successfully triggered workflow runs
/* 163-byte string literal not displayed */
Collisions []*IdempotencyCollisionError `protobuf:"bytes,2,rep,name=collisions,proto3" json:"collisions,omitempty"` // the idempotency collision errors
// contains filtered or unexported fields
}
func (*BulkTriggerIdempotencyCollisionError) Descriptor
deprecated
added in
v0.95.0
func (*BulkTriggerIdempotencyCollisionError) Descriptor() ([]byte, []int)
Deprecated: Use BulkTriggerIdempotencyCollisionError.ProtoReflect.Descriptor instead.
func (*BulkTriggerIdempotencyCollisionError) GetCollisions ¶ added in v0.95.0
func (x *BulkTriggerIdempotencyCollisionError) GetCollisions() []*IdempotencyCollisionError
func (*BulkTriggerIdempotencyCollisionError) GetSuccessfulWorkflowRunExternalIds ¶ added in v0.95.0
func (x *BulkTriggerIdempotencyCollisionError) GetSuccessfulWorkflowRunExternalIds() []string
func (*BulkTriggerIdempotencyCollisionError) ProtoMessage ¶ added in v0.95.0
func (*BulkTriggerIdempotencyCollisionError) ProtoMessage()
func (*BulkTriggerIdempotencyCollisionError) ProtoReflect ¶ added in v0.95.0
func (x *BulkTriggerIdempotencyCollisionError) ProtoReflect() protoreflect.Message
func (*BulkTriggerIdempotencyCollisionError) Reset ¶ added in v0.95.0
func (x *BulkTriggerIdempotencyCollisionError) Reset()
func (*BulkTriggerIdempotencyCollisionError) String ¶ added in v0.95.0
func (x *BulkTriggerIdempotencyCollisionError) String() string
type CancelTasksRequest ¶
type CancelTasksRequest struct {
ExternalIds []string `protobuf:"bytes,1,rep,name=external_ids,json=externalIds,proto3" json:"external_ids,omitempty"` // a list of external UUIDs
Filter *TasksFilter `protobuf:"bytes,2,opt,name=filter,proto3,oneof" json:"filter,omitempty"`
// contains filtered or unexported fields
}
func (*CancelTasksRequest) Descriptor
deprecated
func (*CancelTasksRequest) Descriptor() ([]byte, []int)
Deprecated: Use CancelTasksRequest.ProtoReflect.Descriptor instead.
func (*CancelTasksRequest) GetExternalIds ¶
func (x *CancelTasksRequest) GetExternalIds() []string
func (*CancelTasksRequest) GetFilter ¶
func (x *CancelTasksRequest) GetFilter() *TasksFilter
func (*CancelTasksRequest) ProtoMessage ¶
func (*CancelTasksRequest) ProtoMessage()
func (*CancelTasksRequest) ProtoReflect ¶
func (x *CancelTasksRequest) ProtoReflect() protoreflect.Message
func (*CancelTasksRequest) Reset ¶
func (x *CancelTasksRequest) Reset()
func (*CancelTasksRequest) String ¶
func (x *CancelTasksRequest) String() string
type CancelTasksResponse ¶
type CancelTasksResponse struct {
CancelledTasks []string `protobuf:"bytes,1,rep,name=cancelled_tasks,json=cancelledTasks,proto3" json:"cancelled_tasks,omitempty"`
// contains filtered or unexported fields
}
func (*CancelTasksResponse) Descriptor
deprecated
func (*CancelTasksResponse) Descriptor() ([]byte, []int)
Deprecated: Use CancelTasksResponse.ProtoReflect.Descriptor instead.
func (*CancelTasksResponse) GetCancelledTasks ¶
func (x *CancelTasksResponse) GetCancelledTasks() []string
func (*CancelTasksResponse) ProtoMessage ¶
func (*CancelTasksResponse) ProtoMessage()
func (*CancelTasksResponse) ProtoReflect ¶
func (x *CancelTasksResponse) ProtoReflect() protoreflect.Message
func (*CancelTasksResponse) Reset ¶
func (x *CancelTasksResponse) Reset()
func (*CancelTasksResponse) String ¶
func (x *CancelTasksResponse) String() string
type Concurrency ¶
type Concurrency struct {
Expression string `protobuf:"bytes,1,opt,name=expression,proto3" json:"expression,omitempty"` // (required) the expression to use for concurrency
MaxRuns *int32 `protobuf:"varint,2,opt,name=max_runs,json=maxRuns,proto3,oneof" json:"max_runs,omitempty"` // (optional) the maximum number of concurrent workflow runs, default 1
LimitStrategy *ConcurrencyLimitStrategy `` // (optional) the strategy to use when the concurrency limit is reached, default CANCEL_IN_PROGRESS
/* 140-byte string literal not displayed */
Name *string `protobuf:"bytes,4,opt,name=name,proto3,oneof" json:"name,omitempty"` // (required when is_tenant_scoped) the strategy name; unique per tenant for tenant-scoped strategies
IsTenantScoped *bool `protobuf:"varint,5,opt,name=is_tenant_scoped,json=isTenantScoped,proto3,oneof" json:"is_tenant_scoped,omitempty"` // (optional) when true, the entry is a tenant-scoped strategy shared across workflows, default false
MaxRunsExpression *string `protobuf:"bytes,6,opt,name=max_runs_expression,json=maxRunsExpression,proto3,oneof" json:"max_runs_expression,omitempty"` // (optional) CEL expression over task input returning the max runs for that task's concurrency group; the group's effective limit is the value from its most recently created task. Overrides max_runs per group; a non-integer or negative result fails the task, and 0 holds the group until a newer task raises the limit
// contains filtered or unexported fields
}
Concurrency declares one entry in a concurrency chain. Entries are processed in array order, so a tenant-scoped entry may come before or after a workflow-scoped one.
A tenant-scoped entry (is_tenant_scoped = true) defines (or updates in place) a strategy shared across workflows, keyed by name: every task declaring the same name consumes the same concurrency limit. Registrations whose chains order the same tenant-scoped strategies inconsistently are rejected, since inconsistent orders can deadlock.
func (*Concurrency) Descriptor
deprecated
func (*Concurrency) Descriptor() ([]byte, []int)
Deprecated: Use Concurrency.ProtoReflect.Descriptor instead.
func (*Concurrency) GetExpression ¶
func (x *Concurrency) GetExpression() string
func (*Concurrency) GetIsTenantScoped ¶ added in v0.105.20
func (x *Concurrency) GetIsTenantScoped() bool
func (*Concurrency) GetLimitStrategy ¶
func (x *Concurrency) GetLimitStrategy() ConcurrencyLimitStrategy
func (*Concurrency) GetMaxRuns ¶
func (x *Concurrency) GetMaxRuns() int32
func (*Concurrency) GetMaxRunsExpression ¶ added in v0.105.22
func (x *Concurrency) GetMaxRunsExpression() string
func (*Concurrency) GetName ¶ added in v0.105.20
func (x *Concurrency) GetName() string
func (*Concurrency) ProtoMessage ¶
func (*Concurrency) ProtoMessage()
func (*Concurrency) ProtoReflect ¶
func (x *Concurrency) ProtoReflect() protoreflect.Message
func (*Concurrency) Reset ¶
func (x *Concurrency) Reset()
func (*Concurrency) String ¶
func (x *Concurrency) String() string
type ConcurrencyLimitStrategy ¶
type ConcurrencyLimitStrategy int32
const ( ConcurrencyLimitStrategy_CANCEL_IN_PROGRESS ConcurrencyLimitStrategy = 0 ConcurrencyLimitStrategy_DROP_NEWEST ConcurrencyLimitStrategy = 1 // deprecated ConcurrencyLimitStrategy_QUEUE_NEWEST ConcurrencyLimitStrategy = 2 // deprecated ConcurrencyLimitStrategy_GROUP_ROUND_ROBIN ConcurrencyLimitStrategy = 3 ConcurrencyLimitStrategy_CANCEL_NEWEST ConcurrencyLimitStrategy = 4 ConcurrencyLimitStrategy_CANCEL_QUEUED_EXCEPT_NEWEST ConcurrencyLimitStrategy = 5 ConcurrencyLimitStrategy_CANCEL_QUEUED_EXCEPT_OLDEST ConcurrencyLimitStrategy = 6 )
func (ConcurrencyLimitStrategy) Descriptor ¶
func (ConcurrencyLimitStrategy) Descriptor() protoreflect.EnumDescriptor
func (ConcurrencyLimitStrategy) Enum ¶
func (x ConcurrencyLimitStrategy) Enum() *ConcurrencyLimitStrategy
func (ConcurrencyLimitStrategy) EnumDescriptor
deprecated
func (ConcurrencyLimitStrategy) EnumDescriptor() ([]byte, []int)
Deprecated: Use ConcurrencyLimitStrategy.Descriptor instead.
func (ConcurrencyLimitStrategy) Number ¶
func (x ConcurrencyLimitStrategy) Number() protoreflect.EnumNumber
func (ConcurrencyLimitStrategy) String ¶
func (x ConcurrencyLimitStrategy) String() string
func (ConcurrencyLimitStrategy) Type ¶
func (ConcurrencyLimitStrategy) Type() protoreflect.EnumType
type CreateTaskOpts ¶
type CreateTaskOpts struct {
ReadableId string `protobuf:"bytes,1,opt,name=readable_id,json=readableId,proto3" json:"readable_id,omitempty"` // (required) the task name
Action string `protobuf:"bytes,2,opt,name=action,proto3" json:"action,omitempty"` // (required) the task action id
Timeout string `protobuf:"bytes,3,opt,name=timeout,proto3" json:"timeout,omitempty"` // (optional) the task timeout
Inputs string `protobuf:"bytes,4,opt,name=inputs,proto3" json:"inputs,omitempty"` // (optional) the task inputs, assuming string representation of JSON
Parents []string `protobuf:"bytes,5,rep,name=parents,proto3" json:"parents,omitempty"` // (optional) the task parents. if none are passed in, this is a root task
Retries int32 `protobuf:"varint,6,opt,name=retries,proto3" json:"retries,omitempty"` // (optional) the number of retries for the task, default 0
RateLimits []*CreateTaskRateLimit `protobuf:"bytes,7,rep,name=rate_limits,json=rateLimits,proto3" json:"rate_limits,omitempty"` // (optional) the rate limits for the task
WorkerLabels map[string]*DesiredWorkerLabels `` // (optional) the desired worker affinity state for the task
/* 185-byte string literal not displayed */
BackoffFactor *float32 `protobuf:"fixed32,9,opt,name=backoff_factor,json=backoffFactor,proto3,oneof" json:"backoff_factor,omitempty"` // (optional) the retry backoff factor for the task
BackoffMaxSeconds *int32 `protobuf:"varint,10,opt,name=backoff_max_seconds,json=backoffMaxSeconds,proto3,oneof" json:"backoff_max_seconds,omitempty"` // (optional) the maximum backoff time for the task
Concurrency []*Concurrency `protobuf:"bytes,11,rep,name=concurrency,proto3" json:"concurrency,omitempty"` // (optional) the task concurrency options
Conditions *TaskConditions `protobuf:"bytes,12,opt,name=conditions,proto3,oneof" json:"conditions,omitempty"` // (optional) the task conditions for creating the task
ScheduleTimeout *string `protobuf:"bytes,13,opt,name=schedule_timeout,json=scheduleTimeout,proto3,oneof" json:"schedule_timeout,omitempty"` // (optional) the timeout for the schedule
IsDurable bool `protobuf:"varint,14,opt,name=is_durable,json=isDurable,proto3" json:"is_durable,omitempty"` // (optional) whether the task is durable
SlotRequests map[string]int32 `` // (optional) slot requests (slot_type -> units)
/* 187-byte string literal not displayed */
Batch *TaskBatchConfig `protobuf:"bytes,16,opt,name=batch,proto3,oneof" json:"batch,omitempty"` // (optional) batch execution configuration
// contains filtered or unexported fields
}
CreateTaskOpts represents options to create a task.
func (*CreateTaskOpts) Descriptor
deprecated
func (*CreateTaskOpts) Descriptor() ([]byte, []int)
Deprecated: Use CreateTaskOpts.ProtoReflect.Descriptor instead.
func (*CreateTaskOpts) GetAction ¶
func (x *CreateTaskOpts) GetAction() string
func (*CreateTaskOpts) GetBackoffFactor ¶
func (x *CreateTaskOpts) GetBackoffFactor() float32
func (*CreateTaskOpts) GetBackoffMaxSeconds ¶
func (x *CreateTaskOpts) GetBackoffMaxSeconds() int32
func (*CreateTaskOpts) GetBatch ¶ added in v0.98.0
func (x *CreateTaskOpts) GetBatch() *TaskBatchConfig
func (*CreateTaskOpts) GetConcurrency ¶
func (x *CreateTaskOpts) GetConcurrency() []*Concurrency
func (*CreateTaskOpts) GetConditions ¶
func (x *CreateTaskOpts) GetConditions() *TaskConditions
func (*CreateTaskOpts) GetInputs ¶
func (x *CreateTaskOpts) GetInputs() string
func (*CreateTaskOpts) GetIsDurable ¶ added in v0.78.27
func (x *CreateTaskOpts) GetIsDurable() bool
func (*CreateTaskOpts) GetParents ¶
func (x *CreateTaskOpts) GetParents() []string
func (*CreateTaskOpts) GetRateLimits ¶
func (x *CreateTaskOpts) GetRateLimits() []*CreateTaskRateLimit
func (*CreateTaskOpts) GetReadableId ¶
func (x *CreateTaskOpts) GetReadableId() string
func (*CreateTaskOpts) GetRetries ¶
func (x *CreateTaskOpts) GetRetries() int32
func (*CreateTaskOpts) GetScheduleTimeout ¶
func (x *CreateTaskOpts) GetScheduleTimeout() string
func (*CreateTaskOpts) GetSlotRequests ¶ added in v0.78.27
func (x *CreateTaskOpts) GetSlotRequests() map[string]int32
func (*CreateTaskOpts) GetTimeout ¶
func (x *CreateTaskOpts) GetTimeout() string
func (*CreateTaskOpts) GetWorkerLabels ¶
func (x *CreateTaskOpts) GetWorkerLabels() map[string]*DesiredWorkerLabels
func (*CreateTaskOpts) ProtoMessage ¶
func (*CreateTaskOpts) ProtoMessage()
func (*CreateTaskOpts) ProtoReflect ¶
func (x *CreateTaskOpts) ProtoReflect() protoreflect.Message
func (*CreateTaskOpts) Reset ¶
func (x *CreateTaskOpts) Reset()
func (*CreateTaskOpts) String ¶
func (x *CreateTaskOpts) String() string
type CreateTaskRateLimit ¶
type CreateTaskRateLimit struct {
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"` // (required) the key for the rate limit
Units *int32 `protobuf:"varint,2,opt,name=units,proto3,oneof" json:"units,omitempty"` // (optional) the number of units this task consumes
KeyExpr *string `protobuf:"bytes,3,opt,name=key_expr,json=keyExpr,proto3,oneof" json:"key_expr,omitempty"` // (optional) a CEL expression for determining the rate limit key
UnitsExpr *string `protobuf:"bytes,4,opt,name=units_expr,json=unitsExpr,proto3,oneof" json:"units_expr,omitempty"` // (optional) a CEL expression for determining the number of units consumed
LimitValuesExpr *string `protobuf:"bytes,5,opt,name=limit_values_expr,json=limitValuesExpr,proto3,oneof" json:"limit_values_expr,omitempty"` // (optional) a CEL expression for determining the total amount of rate limit units
Duration *RateLimitDuration `protobuf:"varint,6,opt,name=duration,proto3,enum=v1.RateLimitDuration,oneof" json:"duration,omitempty"` // (optional) the default rate limit window to use for dynamic rate limits
// contains filtered or unexported fields
}
func (*CreateTaskRateLimit) Descriptor
deprecated
func (*CreateTaskRateLimit) Descriptor() ([]byte, []int)
Deprecated: Use CreateTaskRateLimit.ProtoReflect.Descriptor instead.
func (*CreateTaskRateLimit) GetDuration ¶
func (x *CreateTaskRateLimit) GetDuration() RateLimitDuration
func (*CreateTaskRateLimit) GetKey ¶
func (x *CreateTaskRateLimit) GetKey() string
func (*CreateTaskRateLimit) GetKeyExpr ¶
func (x *CreateTaskRateLimit) GetKeyExpr() string
func (*CreateTaskRateLimit) GetLimitValuesExpr ¶
func (x *CreateTaskRateLimit) GetLimitValuesExpr() string
func (*CreateTaskRateLimit) GetUnits ¶
func (x *CreateTaskRateLimit) GetUnits() int32
func (*CreateTaskRateLimit) GetUnitsExpr ¶
func (x *CreateTaskRateLimit) GetUnitsExpr() string
func (*CreateTaskRateLimit) ProtoMessage ¶
func (*CreateTaskRateLimit) ProtoMessage()
func (*CreateTaskRateLimit) ProtoReflect ¶
func (x *CreateTaskRateLimit) ProtoReflect() protoreflect.Message
func (*CreateTaskRateLimit) Reset ¶
func (x *CreateTaskRateLimit) Reset()
func (*CreateTaskRateLimit) String ¶
func (x *CreateTaskRateLimit) String() string
type CreateWorkflowVersionRequest ¶
type CreateWorkflowVersionRequest struct {
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` // (required) the workflow name
Description string `protobuf:"bytes,2,opt,name=description,proto3" json:"description,omitempty"` // (optional) the workflow description
Version string `protobuf:"bytes,3,opt,name=version,proto3" json:"version,omitempty"` // (optional) the workflow version
EventTriggers []string `protobuf:"bytes,4,rep,name=event_triggers,json=eventTriggers,proto3" json:"event_triggers,omitempty"` // (optional) event triggers for the workflow
CronTriggers []string `protobuf:"bytes,5,rep,name=cron_triggers,json=cronTriggers,proto3" json:"cron_triggers,omitempty"` // (optional) cron triggers for the workflow
Tasks []*CreateTaskOpts `protobuf:"bytes,6,rep,name=tasks,proto3" json:"tasks,omitempty"` // (required) the workflow jobs
// Deprecated: use concurrency_arr instead
Concurrency *Concurrency `protobuf:"bytes,7,opt,name=concurrency,proto3" json:"concurrency,omitempty"` // (optional) the workflow concurrency options
CronInput *string `protobuf:"bytes,8,opt,name=cron_input,json=cronInput,proto3,oneof" json:"cron_input,omitempty"` // (optional) the input for the cron trigger
OnFailureTask *CreateTaskOpts `protobuf:"bytes,9,opt,name=on_failure_task,json=onFailureTask,proto3,oneof" json:"on_failure_task,omitempty"` // (optional) the job to run on failure
Sticky *StickyStrategy `protobuf:"varint,10,opt,name=sticky,proto3,enum=v1.StickyStrategy,oneof" json:"sticky,omitempty"` // (optional) the sticky strategy for assigning tasks to workers
DefaultPriority *int32 `protobuf:"varint,11,opt,name=default_priority,json=defaultPriority,proto3,oneof" json:"default_priority,omitempty"` // (optional) the default priority for the workflow
ConcurrencyArr []*Concurrency `protobuf:"bytes,12,rep,name=concurrency_arr,json=concurrencyArr,proto3" json:"concurrency_arr,omitempty"` // (optional) the workflow concurrency options
DefaultFilters []*DefaultFilter `protobuf:"bytes,13,rep,name=default_filters,json=defaultFilters,proto3" json:"default_filters,omitempty"` // (optional) the default filters for the workflow
InputJsonSchema []byte `protobuf:"bytes,14,opt,name=input_json_schema,json=inputJsonSchema,proto3,oneof" json:"input_json_schema,omitempty"` // (optional) the JSON schema for the workflow input
Idempotency *IdempotencyConfig `protobuf:"bytes,15,opt,name=idempotency,proto3,oneof" json:"idempotency,omitempty"` // (optional) idempotency configuration for the workflow
// contains filtered or unexported fields
}
CreateWorkflowVersionRequest represents options to create a workflow version.
func (*CreateWorkflowVersionRequest) Descriptor
deprecated
func (*CreateWorkflowVersionRequest) Descriptor() ([]byte, []int)
Deprecated: Use CreateWorkflowVersionRequest.ProtoReflect.Descriptor instead.
func (*CreateWorkflowVersionRequest) GetConcurrency ¶
func (x *CreateWorkflowVersionRequest) GetConcurrency() *Concurrency
func (*CreateWorkflowVersionRequest) GetConcurrencyArr ¶
func (x *CreateWorkflowVersionRequest) GetConcurrencyArr() []*Concurrency
func (*CreateWorkflowVersionRequest) GetCronInput ¶
func (x *CreateWorkflowVersionRequest) GetCronInput() string
func (*CreateWorkflowVersionRequest) GetCronTriggers ¶
func (x *CreateWorkflowVersionRequest) GetCronTriggers() []string
func (*CreateWorkflowVersionRequest) GetDefaultFilters ¶
func (x *CreateWorkflowVersionRequest) GetDefaultFilters() []*DefaultFilter
func (*CreateWorkflowVersionRequest) GetDefaultPriority ¶
func (x *CreateWorkflowVersionRequest) GetDefaultPriority() int32
func (*CreateWorkflowVersionRequest) GetDescription ¶
func (x *CreateWorkflowVersionRequest) GetDescription() string
func (*CreateWorkflowVersionRequest) GetEventTriggers ¶
func (x *CreateWorkflowVersionRequest) GetEventTriggers() []string
func (*CreateWorkflowVersionRequest) GetIdempotency ¶ added in v0.95.0
func (x *CreateWorkflowVersionRequest) GetIdempotency() *IdempotencyConfig
func (*CreateWorkflowVersionRequest) GetInputJsonSchema ¶ added in v0.77.33
func (x *CreateWorkflowVersionRequest) GetInputJsonSchema() []byte
func (*CreateWorkflowVersionRequest) GetName ¶
func (x *CreateWorkflowVersionRequest) GetName() string
func (*CreateWorkflowVersionRequest) GetOnFailureTask ¶
func (x *CreateWorkflowVersionRequest) GetOnFailureTask() *CreateTaskOpts
func (*CreateWorkflowVersionRequest) GetSticky ¶
func (x *CreateWorkflowVersionRequest) GetSticky() StickyStrategy
func (*CreateWorkflowVersionRequest) GetTasks ¶
func (x *CreateWorkflowVersionRequest) GetTasks() []*CreateTaskOpts
func (*CreateWorkflowVersionRequest) GetVersion ¶
func (x *CreateWorkflowVersionRequest) GetVersion() string
func (*CreateWorkflowVersionRequest) ProtoMessage ¶
func (*CreateWorkflowVersionRequest) ProtoMessage()
func (*CreateWorkflowVersionRequest) ProtoReflect ¶
func (x *CreateWorkflowVersionRequest) ProtoReflect() protoreflect.Message
func (*CreateWorkflowVersionRequest) Reset ¶
func (x *CreateWorkflowVersionRequest) Reset()
func (*CreateWorkflowVersionRequest) String ¶
func (x *CreateWorkflowVersionRequest) String() string
type CreateWorkflowVersionResponse ¶
type CreateWorkflowVersionResponse struct {
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
WorkflowId string `protobuf:"bytes,2,opt,name=workflow_id,json=workflowId,proto3" json:"workflow_id,omitempty"`
// contains filtered or unexported fields
}
CreateWorkflowVersionResponse represents the response after creating a workflow version.
func (*CreateWorkflowVersionResponse) Descriptor
deprecated
func (*CreateWorkflowVersionResponse) Descriptor() ([]byte, []int)
Deprecated: Use CreateWorkflowVersionResponse.ProtoReflect.Descriptor instead.
func (*CreateWorkflowVersionResponse) GetId ¶
func (x *CreateWorkflowVersionResponse) GetId() string
func (*CreateWorkflowVersionResponse) GetWorkflowId ¶
func (x *CreateWorkflowVersionResponse) GetWorkflowId() string
func (*CreateWorkflowVersionResponse) ProtoMessage ¶
func (*CreateWorkflowVersionResponse) ProtoMessage()
func (*CreateWorkflowVersionResponse) ProtoReflect ¶
func (x *CreateWorkflowVersionResponse) ProtoReflect() protoreflect.Message
func (*CreateWorkflowVersionResponse) Reset ¶
func (x *CreateWorkflowVersionResponse) Reset()
func (*CreateWorkflowVersionResponse) String ¶
func (x *CreateWorkflowVersionResponse) String() string
type DefaultFilter ¶
type DefaultFilter struct {
Expression string `protobuf:"bytes,1,opt,name=expression,proto3" json:"expression,omitempty"` // (required) the CEL expression for the filter
Scope string `protobuf:"bytes,2,opt,name=scope,proto3" json:"scope,omitempty"` // (required) the scope of the filter
Payload []byte `protobuf:"bytes,3,opt,name=payload,proto3,oneof" json:"payload,omitempty"` // (optional) the payload for the filter, if any. A JSON object as a string.
// contains filtered or unexported fields
}
func (*DefaultFilter) Descriptor
deprecated
func (*DefaultFilter) Descriptor() ([]byte, []int)
Deprecated: Use DefaultFilter.ProtoReflect.Descriptor instead.
func (*DefaultFilter) GetExpression ¶
func (x *DefaultFilter) GetExpression() string
func (*DefaultFilter) GetPayload ¶
func (x *DefaultFilter) GetPayload() []byte
func (*DefaultFilter) GetScope ¶
func (x *DefaultFilter) GetScope() string
func (*DefaultFilter) ProtoMessage ¶
func (*DefaultFilter) ProtoMessage()
func (*DefaultFilter) ProtoReflect ¶
func (x *DefaultFilter) ProtoReflect() protoreflect.Message
func (*DefaultFilter) Reset ¶
func (x *DefaultFilter) Reset()
func (*DefaultFilter) String ¶
func (x *DefaultFilter) String() string
type DesiredWorkerLabels ¶
type DesiredWorkerLabels struct {
// value of the affinity
StrValue *string `protobuf:"bytes,1,opt,name=str_value,json=strValue,proto3,oneof" json:"str_value,omitempty"`
IntValue *int32 `protobuf:"varint,2,opt,name=int_value,json=intValue,proto3,oneof" json:"int_value,omitempty"`
// *
// (optional) Specifies whether the affinity setting is required.
// If required, the worker will not accept actions that do not have a truthy affinity setting.
//
// Defaults to false.
Required *bool `protobuf:"varint,3,opt,name=required,proto3,oneof" json:"required,omitempty"`
// *
// (optional) Specifies the comparator for the affinity setting.
// If not set, the default is EQUAL.
Comparator *WorkerLabelComparator `protobuf:"varint,4,opt,name=comparator,proto3,enum=v1.WorkerLabelComparator,oneof" json:"comparator,omitempty"`
// *
// (optional) Specifies the weight of the affinity setting.
// If not set, the default is 100.
Weight *int32 `protobuf:"varint,5,opt,name=weight,proto3,oneof" json:"weight,omitempty"`
// contains filtered or unexported fields
}
func (*DesiredWorkerLabels) Descriptor
deprecated
func (*DesiredWorkerLabels) Descriptor() ([]byte, []int)
Deprecated: Use DesiredWorkerLabels.ProtoReflect.Descriptor instead.
func (*DesiredWorkerLabels) GetComparator ¶
func (x *DesiredWorkerLabels) GetComparator() WorkerLabelComparator
func (*DesiredWorkerLabels) GetIntValue ¶
func (x *DesiredWorkerLabels) GetIntValue() int32
func (*DesiredWorkerLabels) GetRequired ¶
func (x *DesiredWorkerLabels) GetRequired() bool
func (*DesiredWorkerLabels) GetStrValue ¶
func (x *DesiredWorkerLabels) GetStrValue() string
func (*DesiredWorkerLabels) GetWeight ¶
func (x *DesiredWorkerLabels) GetWeight() int32
func (*DesiredWorkerLabels) ProtoMessage ¶
func (*DesiredWorkerLabels) ProtoMessage()
func (*DesiredWorkerLabels) ProtoReflect ¶
func (x *DesiredWorkerLabels) ProtoReflect() protoreflect.Message
func (*DesiredWorkerLabels) Reset ¶
func (x *DesiredWorkerLabels) Reset()
func (*DesiredWorkerLabels) String ¶
func (x *DesiredWorkerLabels) String() string
type DurableEvent ¶
type DurableEvent struct {
TaskId string `protobuf:"bytes,1,opt,name=task_id,json=taskId,proto3" json:"task_id,omitempty"`
SignalKey string `protobuf:"bytes,2,opt,name=signal_key,json=signalKey,proto3" json:"signal_key,omitempty"`
Data []byte `protobuf:"bytes,3,opt,name=data,proto3" json:"data,omitempty"` // the data for the event
// contains filtered or unexported fields
}
func (*DurableEvent) Descriptor
deprecated
func (*DurableEvent) Descriptor() ([]byte, []int)
Deprecated: Use DurableEvent.ProtoReflect.Descriptor instead.
func (*DurableEvent) GetData ¶
func (x *DurableEvent) GetData() []byte
func (*DurableEvent) GetSignalKey ¶
func (x *DurableEvent) GetSignalKey() string
func (*DurableEvent) GetTaskId ¶
func (x *DurableEvent) GetTaskId() string
func (*DurableEvent) ProtoMessage ¶
func (*DurableEvent) ProtoMessage()
func (*DurableEvent) ProtoReflect ¶
func (x *DurableEvent) ProtoReflect() protoreflect.Message
func (*DurableEvent) Reset ¶
func (x *DurableEvent) Reset()
func (*DurableEvent) String ¶
func (x *DurableEvent) String() string
type DurableEventListenerConditions ¶
type DurableEventListenerConditions struct {
SleepConditions []*SleepMatchCondition `protobuf:"bytes,1,rep,name=sleep_conditions,json=sleepConditions,proto3" json:"sleep_conditions,omitempty"`
UserEventConditions []*UserEventMatchCondition `protobuf:"bytes,2,rep,name=user_event_conditions,json=userEventConditions,proto3" json:"user_event_conditions,omitempty"`
// contains filtered or unexported fields
}
func (*DurableEventListenerConditions) Descriptor
deprecated
func (*DurableEventListenerConditions) Descriptor() ([]byte, []int)
Deprecated: Use DurableEventListenerConditions.ProtoReflect.Descriptor instead.
func (*DurableEventListenerConditions) GetSleepConditions ¶
func (x *DurableEventListenerConditions) GetSleepConditions() []*SleepMatchCondition
func (*DurableEventListenerConditions) GetUserEventConditions ¶
func (x *DurableEventListenerConditions) GetUserEventConditions() []*UserEventMatchCondition
func (*DurableEventListenerConditions) ProtoMessage ¶
func (*DurableEventListenerConditions) ProtoMessage()
func (*DurableEventListenerConditions) ProtoReflect ¶
func (x *DurableEventListenerConditions) ProtoReflect() protoreflect.Message
func (*DurableEventListenerConditions) Reset ¶
func (x *DurableEventListenerConditions) Reset()
func (*DurableEventListenerConditions) String ¶
func (x *DurableEventListenerConditions) String() string
type DurableEventLogEntryRef ¶ added in v0.80.0
type DurableEventLogEntryRef struct {
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
InvocationCount int32 `protobuf:"varint,2,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
BranchId int64 `protobuf:"varint,3,opt,name=branch_id,json=branchId,proto3" json:"branch_id,omitempty"`
NodeId int64 `protobuf:"varint,4,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"`
// contains filtered or unexported fields
}
func (*DurableEventLogEntryRef) Descriptor
deprecated
added in
v0.80.0
func (*DurableEventLogEntryRef) Descriptor() ([]byte, []int)
Deprecated: Use DurableEventLogEntryRef.ProtoReflect.Descriptor instead.
func (*DurableEventLogEntryRef) GetBranchId ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) GetBranchId() int64
func (*DurableEventLogEntryRef) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) GetDurableTaskExternalId() string
func (*DurableEventLogEntryRef) GetInvocationCount ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) GetInvocationCount() int32
func (*DurableEventLogEntryRef) GetNodeId ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) GetNodeId() int64
func (*DurableEventLogEntryRef) ProtoMessage ¶ added in v0.80.0
func (*DurableEventLogEntryRef) ProtoMessage()
func (*DurableEventLogEntryRef) ProtoReflect ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) ProtoReflect() protoreflect.Message
func (*DurableEventLogEntryRef) Reset ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) Reset()
func (*DurableEventLogEntryRef) String ¶ added in v0.80.0
func (x *DurableEventLogEntryRef) String() string
type DurableTaskAwaitedCompletedEntry ¶ added in v0.80.0
type DurableTaskAwaitedCompletedEntry struct {
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
BranchId int64 `protobuf:"varint,2,opt,name=branch_id,json=branchId,proto3" json:"branch_id,omitempty"`
NodeId int64 `protobuf:"varint,3,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"`
InvocationCount int32 `protobuf:"varint,4,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskAwaitedCompletedEntry) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskAwaitedCompletedEntry) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskAwaitedCompletedEntry.ProtoReflect.Descriptor instead.
func (*DurableTaskAwaitedCompletedEntry) GetBranchId ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) GetBranchId() int64
func (*DurableTaskAwaitedCompletedEntry) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) GetDurableTaskExternalId() string
func (*DurableTaskAwaitedCompletedEntry) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) GetInvocationCount() int32
func (*DurableTaskAwaitedCompletedEntry) GetNodeId ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) GetNodeId() int64
func (*DurableTaskAwaitedCompletedEntry) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskAwaitedCompletedEntry) ProtoMessage()
func (*DurableTaskAwaitedCompletedEntry) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) ProtoReflect() protoreflect.Message
func (*DurableTaskAwaitedCompletedEntry) Reset ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) Reset()
func (*DurableTaskAwaitedCompletedEntry) String ¶ added in v0.80.0
func (x *DurableTaskAwaitedCompletedEntry) String() string
type DurableTaskCompleteMemoRequest ¶ added in v0.80.0
type DurableTaskCompleteMemoRequest struct {
Ref *DurableEventLogEntryRef `protobuf:"bytes,1,opt,name=ref,proto3" json:"ref,omitempty"`
Payload []byte `protobuf:"bytes,2,opt,name=payload,proto3" json:"payload,omitempty"`
MemoKey []byte `protobuf:"bytes,3,opt,name=memo_key,json=memoKey,proto3" json:"memo_key,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskCompleteMemoRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskCompleteMemoRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskCompleteMemoRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskCompleteMemoRequest) GetMemoKey ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) GetMemoKey() []byte
func (*DurableTaskCompleteMemoRequest) GetPayload ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) GetPayload() []byte
func (*DurableTaskCompleteMemoRequest) GetRef ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) GetRef() *DurableEventLogEntryRef
func (*DurableTaskCompleteMemoRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskCompleteMemoRequest) ProtoMessage()
func (*DurableTaskCompleteMemoRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskCompleteMemoRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) Reset()
func (*DurableTaskCompleteMemoRequest) String ¶ added in v0.80.0
func (x *DurableTaskCompleteMemoRequest) String() string
type DurableTaskErrorResponse ¶ added in v0.80.0
type DurableTaskErrorResponse struct {
Ref *DurableEventLogEntryRef `protobuf:"bytes,1,opt,name=ref,proto3" json:"ref,omitempty"`
ErrorType DurableTaskErrorType `protobuf:"varint,2,opt,name=error_type,json=errorType,proto3,enum=v1.DurableTaskErrorType" json:"error_type,omitempty"`
ErrorMessage string `protobuf:"bytes,3,opt,name=error_message,json=errorMessage,proto3" json:"error_message,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskErrorResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskErrorResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskErrorResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskErrorResponse) GetErrorMessage ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) GetErrorMessage() string
func (*DurableTaskErrorResponse) GetErrorType ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) GetErrorType() DurableTaskErrorType
func (*DurableTaskErrorResponse) GetRef ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) GetRef() *DurableEventLogEntryRef
func (*DurableTaskErrorResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskErrorResponse) ProtoMessage()
func (*DurableTaskErrorResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskErrorResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) Reset()
func (*DurableTaskErrorResponse) String ¶ added in v0.80.0
func (x *DurableTaskErrorResponse) String() string
type DurableTaskErrorType ¶ added in v0.80.0
type DurableTaskErrorType int32
const ( DurableTaskErrorType_DURABLE_TASK_ERROR_TYPE_UNSPECIFIED DurableTaskErrorType = 0 DurableTaskErrorType_DURABLE_TASK_ERROR_TYPE_NONDETERMINISM DurableTaskErrorType = 1 )
func (DurableTaskErrorType) Descriptor ¶ added in v0.80.0
func (DurableTaskErrorType) Descriptor() protoreflect.EnumDescriptor
func (DurableTaskErrorType) Enum ¶ added in v0.80.0
func (x DurableTaskErrorType) Enum() *DurableTaskErrorType
func (DurableTaskErrorType) EnumDescriptor
deprecated
added in
v0.80.0
func (DurableTaskErrorType) EnumDescriptor() ([]byte, []int)
Deprecated: Use DurableTaskErrorType.Descriptor instead.
func (DurableTaskErrorType) Number ¶ added in v0.80.0
func (x DurableTaskErrorType) Number() protoreflect.EnumNumber
func (DurableTaskErrorType) String ¶ added in v0.80.0
func (x DurableTaskErrorType) String() string
func (DurableTaskErrorType) Type ¶ added in v0.80.0
func (DurableTaskErrorType) Type() protoreflect.EnumType
type DurableTaskEventLogEntryCompletedResponse ¶ added in v0.80.0
type DurableTaskEventLogEntryCompletedResponse struct {
Ref *DurableEventLogEntryRef `protobuf:"bytes,1,opt,name=ref,proto3" json:"ref,omitempty"`
Payload []byte `protobuf:"bytes,2,opt,name=payload,proto3" json:"payload,omitempty"`
IsFailure bool `protobuf:"varint,3,opt,name=is_failure,json=isFailure,proto3" json:"is_failure,omitempty"`
ErrorMessage *string `protobuf:"bytes,4,opt,name=error_message,json=errorMessage,proto3,oneof" json:"error_message,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskEventLogEntryCompletedResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEventLogEntryCompletedResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEventLogEntryCompletedResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskEventLogEntryCompletedResponse) GetErrorMessage ¶ added in v0.89.3
func (x *DurableTaskEventLogEntryCompletedResponse) GetErrorMessage() string
func (*DurableTaskEventLogEntryCompletedResponse) GetIsFailure ¶ added in v0.89.3
func (x *DurableTaskEventLogEntryCompletedResponse) GetIsFailure() bool
func (*DurableTaskEventLogEntryCompletedResponse) GetPayload ¶ added in v0.80.0
func (x *DurableTaskEventLogEntryCompletedResponse) GetPayload() []byte
func (*DurableTaskEventLogEntryCompletedResponse) GetRef ¶ added in v0.80.0
func (x *DurableTaskEventLogEntryCompletedResponse) GetRef() *DurableEventLogEntryRef
func (*DurableTaskEventLogEntryCompletedResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEventLogEntryCompletedResponse) ProtoMessage()
func (*DurableTaskEventLogEntryCompletedResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEventLogEntryCompletedResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskEventLogEntryCompletedResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskEventLogEntryCompletedResponse) Reset()
func (*DurableTaskEventLogEntryCompletedResponse) String ¶ added in v0.80.0
func (x *DurableTaskEventLogEntryCompletedResponse) String() string
type DurableTaskEventMemoAckResponse ¶ added in v0.80.0
type DurableTaskEventMemoAckResponse struct {
Ref *DurableEventLogEntryRef `protobuf:"bytes,1,opt,name=ref,proto3" json:"ref,omitempty"`
MemoAlreadyExisted bool `protobuf:"varint,2,opt,name=memo_already_existed,json=memoAlreadyExisted,proto3" json:"memo_already_existed,omitempty"`
MemoResultPayload []byte `protobuf:"bytes,3,opt,name=memo_result_payload,json=memoResultPayload,proto3,oneof" json:"memo_result_payload,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskEventMemoAckResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEventMemoAckResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEventMemoAckResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskEventMemoAckResponse) GetMemoAlreadyExisted ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) GetMemoAlreadyExisted() bool
func (*DurableTaskEventMemoAckResponse) GetMemoResultPayload ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) GetMemoResultPayload() []byte
func (*DurableTaskEventMemoAckResponse) GetRef ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) GetRef() *DurableEventLogEntryRef
func (*DurableTaskEventMemoAckResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEventMemoAckResponse) ProtoMessage()
func (*DurableTaskEventMemoAckResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskEventMemoAckResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) Reset()
func (*DurableTaskEventMemoAckResponse) String ¶ added in v0.80.0
func (x *DurableTaskEventMemoAckResponse) String() string
type DurableTaskEventTriggerRunsAckResponse ¶ added in v0.80.0
type DurableTaskEventTriggerRunsAckResponse struct {
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
InvocationCount int32 `protobuf:"varint,2,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
RunEntries []*DurableTaskRunAckEntry `protobuf:"bytes,3,rep,name=run_entries,json=runEntries,proto3" json:"run_entries,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskEventTriggerRunsAckResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEventTriggerRunsAckResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEventTriggerRunsAckResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskEventTriggerRunsAckResponse) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) GetDurableTaskExternalId() string
func (*DurableTaskEventTriggerRunsAckResponse) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) GetInvocationCount() int32
func (*DurableTaskEventTriggerRunsAckResponse) GetRunEntries ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) GetRunEntries() []*DurableTaskRunAckEntry
func (*DurableTaskEventTriggerRunsAckResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEventTriggerRunsAckResponse) ProtoMessage()
func (*DurableTaskEventTriggerRunsAckResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskEventTriggerRunsAckResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) Reset()
func (*DurableTaskEventTriggerRunsAckResponse) String ¶ added in v0.80.0
func (x *DurableTaskEventTriggerRunsAckResponse) String() string
type DurableTaskEventWaitForAckResponse ¶ added in v0.80.0
type DurableTaskEventWaitForAckResponse struct {
Ref *DurableEventLogEntryRef `protobuf:"bytes,1,opt,name=ref,proto3" json:"ref,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskEventWaitForAckResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEventWaitForAckResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEventWaitForAckResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskEventWaitForAckResponse) GetRef ¶ added in v0.80.0
func (x *DurableTaskEventWaitForAckResponse) GetRef() *DurableEventLogEntryRef
func (*DurableTaskEventWaitForAckResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEventWaitForAckResponse) ProtoMessage()
func (*DurableTaskEventWaitForAckResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEventWaitForAckResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskEventWaitForAckResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskEventWaitForAckResponse) Reset()
func (*DurableTaskEventWaitForAckResponse) String ¶ added in v0.80.0
func (x *DurableTaskEventWaitForAckResponse) String() string
type DurableTaskEvictInvocationRequest ¶ added in v0.80.0
type DurableTaskEvictInvocationRequest struct {
InvocationCount int32 `protobuf:"varint,1,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
Reason *string `protobuf:"bytes,3,opt,name=reason,proto3,oneof" json:"reason,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskEvictInvocationRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEvictInvocationRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEvictInvocationRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskEvictInvocationRequest) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) GetDurableTaskExternalId() string
func (*DurableTaskEvictInvocationRequest) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) GetInvocationCount() int32
func (*DurableTaskEvictInvocationRequest) GetReason ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) GetReason() string
func (*DurableTaskEvictInvocationRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEvictInvocationRequest) ProtoMessage()
func (*DurableTaskEvictInvocationRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskEvictInvocationRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) Reset()
func (*DurableTaskEvictInvocationRequest) String ¶ added in v0.80.0
func (x *DurableTaskEvictInvocationRequest) String() string
type DurableTaskEvictionAckResponse ¶ added in v0.80.0
type DurableTaskEvictionAckResponse struct {
InvocationCount int32 `protobuf:"varint,1,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
// contains filtered or unexported fields
}
Sent by the server after recording eviction for an evict_invocation request.
func (*DurableTaskEvictionAckResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskEvictionAckResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskEvictionAckResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskEvictionAckResponse) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskEvictionAckResponse) GetDurableTaskExternalId() string
func (*DurableTaskEvictionAckResponse) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskEvictionAckResponse) GetInvocationCount() int32
func (*DurableTaskEvictionAckResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskEvictionAckResponse) ProtoMessage()
func (*DurableTaskEvictionAckResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskEvictionAckResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskEvictionAckResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskEvictionAckResponse) Reset()
func (*DurableTaskEvictionAckResponse) String ¶ added in v0.80.0
func (x *DurableTaskEvictionAckResponse) String() string
type DurableTaskMemoRequest ¶ added in v0.80.0
type DurableTaskMemoRequest struct {
// The invocation_count is a monotonically increasing count that uniquely identifies an "attempt"
// at running a durable task. Each time the task is started, it gets a new invocation count (which has)
// incremented by one since the previous invocation. This allows the server (and the worker) to have a way of
// differentiating between different attempts of the same task running in different places, to prevent race conditions
// and other problems from duplication. It also allows for older invocations to be evicted cleanly
InvocationCount int32 `protobuf:"varint,1,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
Key []byte `protobuf:"bytes,3,opt,name=key,proto3" json:"key,omitempty"`
// optional payload because we can send a memo request to check if a memo already exists
Payload []byte `protobuf:"bytes,4,opt,name=payload,proto3,oneof" json:"payload,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskMemoRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskMemoRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskMemoRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskMemoRequest) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) GetDurableTaskExternalId() string
func (*DurableTaskMemoRequest) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) GetInvocationCount() int32
func (*DurableTaskMemoRequest) GetKey ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) GetKey() []byte
func (*DurableTaskMemoRequest) GetPayload ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) GetPayload() []byte
func (*DurableTaskMemoRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskMemoRequest) ProtoMessage()
func (*DurableTaskMemoRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskMemoRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) Reset()
func (*DurableTaskMemoRequest) String ¶ added in v0.80.0
func (x *DurableTaskMemoRequest) String() string
type DurableTaskRequest ¶ added in v0.80.0
type DurableTaskRequest struct {
// Types that are assignable to Message:
//
// *DurableTaskRequest_RegisterWorker
// *DurableTaskRequest_Memo
// *DurableTaskRequest_TriggerRuns
// *DurableTaskRequest_WaitFor
// *DurableTaskRequest_EvictInvocation
// *DurableTaskRequest_WorkerStatus
// *DurableTaskRequest_CompleteMemo
Message isDurableTaskRequest_Message `protobuf_oneof:"message"`
// contains filtered or unexported fields
}
func (*DurableTaskRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskRequest) GetCompleteMemo ¶ added in v0.80.0
func (x *DurableTaskRequest) GetCompleteMemo() *DurableTaskCompleteMemoRequest
func (*DurableTaskRequest) GetEvictInvocation ¶ added in v0.80.0
func (x *DurableTaskRequest) GetEvictInvocation() *DurableTaskEvictInvocationRequest
func (*DurableTaskRequest) GetMemo ¶ added in v0.80.0
func (x *DurableTaskRequest) GetMemo() *DurableTaskMemoRequest
func (*DurableTaskRequest) GetMessage ¶ added in v0.80.0
func (m *DurableTaskRequest) GetMessage() isDurableTaskRequest_Message
func (*DurableTaskRequest) GetRegisterWorker ¶ added in v0.80.0
func (x *DurableTaskRequest) GetRegisterWorker() *DurableTaskRequestRegisterWorker
func (*DurableTaskRequest) GetTriggerRuns ¶ added in v0.80.0
func (x *DurableTaskRequest) GetTriggerRuns() *DurableTaskTriggerRunsRequest
func (*DurableTaskRequest) GetWaitFor ¶ added in v0.80.0
func (x *DurableTaskRequest) GetWaitFor() *DurableTaskWaitForRequest
func (*DurableTaskRequest) GetWorkerStatus ¶ added in v0.80.0
func (x *DurableTaskRequest) GetWorkerStatus() *DurableTaskWorkerStatusRequest
func (*DurableTaskRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskRequest) ProtoMessage()
func (*DurableTaskRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskRequest) Reset()
func (*DurableTaskRequest) String ¶ added in v0.80.0
func (x *DurableTaskRequest) String() string
type DurableTaskRequestRegisterWorker ¶ added in v0.80.0
type DurableTaskRequestRegisterWorker struct {
WorkerId string `protobuf:"bytes,1,opt,name=worker_id,json=workerId,proto3" json:"worker_id,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskRequestRegisterWorker) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskRequestRegisterWorker) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskRequestRegisterWorker.ProtoReflect.Descriptor instead.
func (*DurableTaskRequestRegisterWorker) GetWorkerId ¶ added in v0.80.0
func (x *DurableTaskRequestRegisterWorker) GetWorkerId() string
func (*DurableTaskRequestRegisterWorker) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskRequestRegisterWorker) ProtoMessage()
func (*DurableTaskRequestRegisterWorker) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskRequestRegisterWorker) ProtoReflect() protoreflect.Message
func (*DurableTaskRequestRegisterWorker) Reset ¶ added in v0.80.0
func (x *DurableTaskRequestRegisterWorker) Reset()
func (*DurableTaskRequestRegisterWorker) String ¶ added in v0.80.0
func (x *DurableTaskRequestRegisterWorker) String() string
type DurableTaskRequest_CompleteMemo ¶ added in v0.80.0
type DurableTaskRequest_CompleteMemo struct {
CompleteMemo *DurableTaskCompleteMemoRequest `protobuf:"bytes,7,opt,name=complete_memo,json=completeMemo,proto3,oneof"`
}
type DurableTaskRequest_EvictInvocation ¶ added in v0.80.0
type DurableTaskRequest_EvictInvocation struct {
EvictInvocation *DurableTaskEvictInvocationRequest `protobuf:"bytes,5,opt,name=evict_invocation,json=evictInvocation,proto3,oneof"`
}
type DurableTaskRequest_Memo ¶ added in v0.80.0
type DurableTaskRequest_Memo struct {
Memo *DurableTaskMemoRequest `protobuf:"bytes,2,opt,name=memo,proto3,oneof"`
}
type DurableTaskRequest_RegisterWorker ¶ added in v0.80.0
type DurableTaskRequest_RegisterWorker struct {
RegisterWorker *DurableTaskRequestRegisterWorker `protobuf:"bytes,1,opt,name=register_worker,json=registerWorker,proto3,oneof"`
}
type DurableTaskRequest_TriggerRuns ¶ added in v0.80.0
type DurableTaskRequest_TriggerRuns struct {
TriggerRuns *DurableTaskTriggerRunsRequest `protobuf:"bytes,3,opt,name=trigger_runs,json=triggerRuns,proto3,oneof"`
}
type DurableTaskRequest_WaitFor ¶ added in v0.80.0
type DurableTaskRequest_WaitFor struct {
WaitFor *DurableTaskWaitForRequest `protobuf:"bytes,4,opt,name=wait_for,json=waitFor,proto3,oneof"`
}
type DurableTaskRequest_WorkerStatus ¶ added in v0.80.0
type DurableTaskRequest_WorkerStatus struct {
WorkerStatus *DurableTaskWorkerStatusRequest `protobuf:"bytes,6,opt,name=worker_status,json=workerStatus,proto3,oneof"`
}
type DurableTaskResponse ¶ added in v0.80.0
type DurableTaskResponse struct {
// Types that are assignable to Message:
//
// *DurableTaskResponse_RegisterWorker
// *DurableTaskResponse_MemoAck
// *DurableTaskResponse_TriggerRunsAck
// *DurableTaskResponse_WaitForAck
// *DurableTaskResponse_EntryCompleted
// *DurableTaskResponse_Error
// *DurableTaskResponse_EvictionAck
// *DurableTaskResponse_ServerEvict
Message isDurableTaskResponse_Message `protobuf_oneof:"message"`
// contains filtered or unexported fields
}
func (*DurableTaskResponse) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskResponse) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskResponse.ProtoReflect.Descriptor instead.
func (*DurableTaskResponse) GetEntryCompleted ¶ added in v0.80.0
func (x *DurableTaskResponse) GetEntryCompleted() *DurableTaskEventLogEntryCompletedResponse
func (*DurableTaskResponse) GetError ¶ added in v0.80.0
func (x *DurableTaskResponse) GetError() *DurableTaskErrorResponse
func (*DurableTaskResponse) GetEvictionAck ¶ added in v0.80.0
func (x *DurableTaskResponse) GetEvictionAck() *DurableTaskEvictionAckResponse
func (*DurableTaskResponse) GetMemoAck ¶ added in v0.80.0
func (x *DurableTaskResponse) GetMemoAck() *DurableTaskEventMemoAckResponse
func (*DurableTaskResponse) GetMessage ¶ added in v0.80.0
func (m *DurableTaskResponse) GetMessage() isDurableTaskResponse_Message
func (*DurableTaskResponse) GetRegisterWorker ¶ added in v0.80.0
func (x *DurableTaskResponse) GetRegisterWorker() *DurableTaskResponseRegisterWorker
func (*DurableTaskResponse) GetServerEvict ¶ added in v0.80.0
func (x *DurableTaskResponse) GetServerEvict() *DurableTaskServerEvictNotice
func (*DurableTaskResponse) GetTriggerRunsAck ¶ added in v0.80.0
func (x *DurableTaskResponse) GetTriggerRunsAck() *DurableTaskEventTriggerRunsAckResponse
func (*DurableTaskResponse) GetWaitForAck ¶ added in v0.80.0
func (x *DurableTaskResponse) GetWaitForAck() *DurableTaskEventWaitForAckResponse
func (*DurableTaskResponse) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskResponse) ProtoMessage()
func (*DurableTaskResponse) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskResponse) ProtoReflect() protoreflect.Message
func (*DurableTaskResponse) Reset ¶ added in v0.80.0
func (x *DurableTaskResponse) Reset()
func (*DurableTaskResponse) String ¶ added in v0.80.0
func (x *DurableTaskResponse) String() string
type DurableTaskResponseRegisterWorker ¶ added in v0.80.0
type DurableTaskResponseRegisterWorker struct {
WorkerId string `protobuf:"bytes,1,opt,name=worker_id,json=workerId,proto3" json:"worker_id,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskResponseRegisterWorker) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskResponseRegisterWorker) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskResponseRegisterWorker.ProtoReflect.Descriptor instead.
func (*DurableTaskResponseRegisterWorker) GetWorkerId ¶ added in v0.80.0
func (x *DurableTaskResponseRegisterWorker) GetWorkerId() string
func (*DurableTaskResponseRegisterWorker) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskResponseRegisterWorker) ProtoMessage()
func (*DurableTaskResponseRegisterWorker) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskResponseRegisterWorker) ProtoReflect() protoreflect.Message
func (*DurableTaskResponseRegisterWorker) Reset ¶ added in v0.80.0
func (x *DurableTaskResponseRegisterWorker) Reset()
func (*DurableTaskResponseRegisterWorker) String ¶ added in v0.80.0
func (x *DurableTaskResponseRegisterWorker) String() string
type DurableTaskResponse_EntryCompleted ¶ added in v0.80.0
type DurableTaskResponse_EntryCompleted struct {
EntryCompleted *DurableTaskEventLogEntryCompletedResponse `protobuf:"bytes,5,opt,name=entry_completed,json=entryCompleted,proto3,oneof"`
}
type DurableTaskResponse_Error ¶ added in v0.80.0
type DurableTaskResponse_Error struct {
Error *DurableTaskErrorResponse `protobuf:"bytes,6,opt,name=error,proto3,oneof"`
}
type DurableTaskResponse_EvictionAck ¶ added in v0.80.0
type DurableTaskResponse_EvictionAck struct {
EvictionAck *DurableTaskEvictionAckResponse `protobuf:"bytes,7,opt,name=eviction_ack,json=evictionAck,proto3,oneof"`
}
type DurableTaskResponse_MemoAck ¶ added in v0.80.0
type DurableTaskResponse_MemoAck struct {
MemoAck *DurableTaskEventMemoAckResponse `protobuf:"bytes,2,opt,name=memo_ack,json=memoAck,proto3,oneof"`
}
type DurableTaskResponse_RegisterWorker ¶ added in v0.80.0
type DurableTaskResponse_RegisterWorker struct {
RegisterWorker *DurableTaskResponseRegisterWorker `protobuf:"bytes,1,opt,name=register_worker,json=registerWorker,proto3,oneof"`
}
type DurableTaskResponse_ServerEvict ¶ added in v0.80.0
type DurableTaskResponse_ServerEvict struct {
ServerEvict *DurableTaskServerEvictNotice `protobuf:"bytes,8,opt,name=server_evict,json=serverEvict,proto3,oneof"`
}
type DurableTaskResponse_TriggerRunsAck ¶ added in v0.80.0
type DurableTaskResponse_TriggerRunsAck struct {
TriggerRunsAck *DurableTaskEventTriggerRunsAckResponse `protobuf:"bytes,3,opt,name=trigger_runs_ack,json=triggerRunsAck,proto3,oneof"`
}
type DurableTaskResponse_WaitForAck ¶ added in v0.80.0
type DurableTaskResponse_WaitForAck struct {
WaitForAck *DurableTaskEventWaitForAckResponse `protobuf:"bytes,4,opt,name=wait_for_ack,json=waitForAck,proto3,oneof"`
}
type DurableTaskRunAckEntry ¶ added in v0.80.0
type DurableTaskRunAckEntry struct {
NodeId int64 `protobuf:"varint,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"`
BranchId int64 `protobuf:"varint,2,opt,name=branch_id,json=branchId,proto3" json:"branch_id,omitempty"`
WorkflowRunExternalId string `` /* 128-byte string literal not displayed */
// contains filtered or unexported fields
}
func (*DurableTaskRunAckEntry) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskRunAckEntry) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskRunAckEntry.ProtoReflect.Descriptor instead.
func (*DurableTaskRunAckEntry) GetBranchId ¶ added in v0.80.0
func (x *DurableTaskRunAckEntry) GetBranchId() int64
func (*DurableTaskRunAckEntry) GetNodeId ¶ added in v0.80.0
func (x *DurableTaskRunAckEntry) GetNodeId() int64
func (*DurableTaskRunAckEntry) GetWorkflowRunExternalId ¶ added in v0.80.4
func (x *DurableTaskRunAckEntry) GetWorkflowRunExternalId() string
func (*DurableTaskRunAckEntry) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskRunAckEntry) ProtoMessage()
func (*DurableTaskRunAckEntry) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskRunAckEntry) ProtoReflect() protoreflect.Message
func (*DurableTaskRunAckEntry) Reset ¶ added in v0.80.0
func (x *DurableTaskRunAckEntry) Reset()
func (*DurableTaskRunAckEntry) String ¶ added in v0.80.0
func (x *DurableTaskRunAckEntry) String() string
type DurableTaskServerEvictNotice ¶ added in v0.80.0
type DurableTaskServerEvictNotice struct {
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
InvocationCount int32 `protobuf:"varint,2,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
Reason string `protobuf:"bytes,3,opt,name=reason,proto3" json:"reason,omitempty"`
// contains filtered or unexported fields
}
Sent by the server to notify a worker that its invocation is stale and should be cancelled.
func (*DurableTaskServerEvictNotice) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskServerEvictNotice) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskServerEvictNotice.ProtoReflect.Descriptor instead.
func (*DurableTaskServerEvictNotice) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) GetDurableTaskExternalId() string
func (*DurableTaskServerEvictNotice) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) GetInvocationCount() int32
func (*DurableTaskServerEvictNotice) GetReason ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) GetReason() string
func (*DurableTaskServerEvictNotice) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskServerEvictNotice) ProtoMessage()
func (*DurableTaskServerEvictNotice) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) ProtoReflect() protoreflect.Message
func (*DurableTaskServerEvictNotice) Reset ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) Reset()
func (*DurableTaskServerEvictNotice) String ¶ added in v0.80.0
func (x *DurableTaskServerEvictNotice) String() string
type DurableTaskTriggerRunsRequest ¶ added in v0.80.0
type DurableTaskTriggerRunsRequest struct {
// The invocation_count is a monotonically increasing count that uniquely identifies an "attempt"
// at running a durable task. Each time the task is started, it gets a new invocation count (which has)
// incremented by one since the previous invocation. This allows the server (and the worker) to have a way of
// differentiating between different attempts of the same task running in different places, to prevent race conditions
// and other problems from duplication. It also allows for older invocations to be evicted cleanly
InvocationCount int32 `protobuf:"varint,1,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
TriggerOpts []*TriggerWorkflowRequest `protobuf:"bytes,3,rep,name=trigger_opts,json=triggerOpts,proto3" json:"trigger_opts,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskTriggerRunsRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskTriggerRunsRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskTriggerRunsRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskTriggerRunsRequest) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) GetDurableTaskExternalId() string
func (*DurableTaskTriggerRunsRequest) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) GetInvocationCount() int32
func (*DurableTaskTriggerRunsRequest) GetTriggerOpts ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) GetTriggerOpts() []*TriggerWorkflowRequest
func (*DurableTaskTriggerRunsRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskTriggerRunsRequest) ProtoMessage()
func (*DurableTaskTriggerRunsRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskTriggerRunsRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) Reset()
func (*DurableTaskTriggerRunsRequest) String ¶ added in v0.80.0
func (x *DurableTaskTriggerRunsRequest) String() string
type DurableTaskWaitForRequest ¶ added in v0.80.0
type DurableTaskWaitForRequest struct {
// The invocation_count is a monotonically increasing count that uniquely identifies an "attempt"
// at running a durable task. Each time the task is started, it gets a new invocation count (which has)
// incremented by one since the previous invocation. This allows the server (and the worker) to have a way of
// differentiating between different attempts of the same task running in different places, to prevent race conditions
// and other problems from duplication. It also allows for older invocations to be evicted cleanly
InvocationCount int32 `protobuf:"varint,1,opt,name=invocation_count,json=invocationCount,proto3" json:"invocation_count,omitempty"`
DurableTaskExternalId string `` /* 128-byte string literal not displayed */
// Fields for DURABLE_TASK_TRIGGER_KIND_WAIT_FOR
WaitForConditions *DurableEventListenerConditions `protobuf:"bytes,3,opt,name=wait_for_conditions,json=waitForConditions,proto3,oneof" json:"wait_for_conditions,omitempty"`
// An optional human-readable label for this wait, displayed in the dashboard.
// Example: "Waiting for payment confirmation"
Label *string `protobuf:"bytes,4,opt,name=label,proto3,oneof" json:"label,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskWaitForRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskWaitForRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskWaitForRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskWaitForRequest) GetDurableTaskExternalId ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) GetDurableTaskExternalId() string
func (*DurableTaskWaitForRequest) GetInvocationCount ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) GetInvocationCount() int32
func (*DurableTaskWaitForRequest) GetLabel ¶ added in v0.83.32
func (x *DurableTaskWaitForRequest) GetLabel() string
func (*DurableTaskWaitForRequest) GetWaitForConditions ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) GetWaitForConditions() *DurableEventListenerConditions
func (*DurableTaskWaitForRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskWaitForRequest) ProtoMessage()
func (*DurableTaskWaitForRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskWaitForRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) Reset()
func (*DurableTaskWaitForRequest) String ¶ added in v0.80.0
func (x *DurableTaskWaitForRequest) String() string
type DurableTaskWorkerStatusRequest ¶ added in v0.80.0
type DurableTaskWorkerStatusRequest struct {
WorkerId string `protobuf:"bytes,1,opt,name=worker_id,json=workerId,proto3" json:"worker_id,omitempty"`
WaitingEntries []*DurableTaskAwaitedCompletedEntry `protobuf:"bytes,2,rep,name=waiting_entries,json=waitingEntries,proto3" json:"waiting_entries,omitempty"`
// contains filtered or unexported fields
}
func (*DurableTaskWorkerStatusRequest) Descriptor
deprecated
added in
v0.80.0
func (*DurableTaskWorkerStatusRequest) Descriptor() ([]byte, []int)
Deprecated: Use DurableTaskWorkerStatusRequest.ProtoReflect.Descriptor instead.
func (*DurableTaskWorkerStatusRequest) GetWaitingEntries ¶ added in v0.80.0
func (x *DurableTaskWorkerStatusRequest) GetWaitingEntries() []*DurableTaskAwaitedCompletedEntry
func (*DurableTaskWorkerStatusRequest) GetWorkerId ¶ added in v0.80.0
func (x *DurableTaskWorkerStatusRequest) GetWorkerId() string
func (*DurableTaskWorkerStatusRequest) ProtoMessage ¶ added in v0.80.0
func (*DurableTaskWorkerStatusRequest) ProtoMessage()
func (*DurableTaskWorkerStatusRequest) ProtoReflect ¶ added in v0.80.0
func (x *DurableTaskWorkerStatusRequest) ProtoReflect() protoreflect.Message
func (*DurableTaskWorkerStatusRequest) Reset ¶ added in v0.80.0
func (x *DurableTaskWorkerStatusRequest) Reset()
func (*DurableTaskWorkerStatusRequest) String ¶ added in v0.80.0
func (x *DurableTaskWorkerStatusRequest) String() string
type GetRunDetailsRequest ¶ added in v0.74.9
type GetRunDetailsRequest struct {
ExternalId string `protobuf:"bytes,1,opt,name=external_id,json=externalId,proto3" json:"external_id,omitempty"` // (required) the external id (uuid) of the workflow run
// contains filtered or unexported fields
}
func (*GetRunDetailsRequest) Descriptor
deprecated
added in
v0.74.9
func (*GetRunDetailsRequest) Descriptor() ([]byte, []int)
Deprecated: Use GetRunDetailsRequest.ProtoReflect.Descriptor instead.
func (*GetRunDetailsRequest) GetExternalId ¶ added in v0.74.9
func (x *GetRunDetailsRequest) GetExternalId() string
func (*GetRunDetailsRequest) ProtoMessage ¶ added in v0.74.9
func (*GetRunDetailsRequest) ProtoMessage()
func (*GetRunDetailsRequest) ProtoReflect ¶ added in v0.74.9
func (x *GetRunDetailsRequest) ProtoReflect() protoreflect.Message
func (*GetRunDetailsRequest) Reset ¶ added in v0.74.9
func (x *GetRunDetailsRequest) Reset()
func (*GetRunDetailsRequest) String ¶ added in v0.74.9
func (x *GetRunDetailsRequest) String() string
type GetRunDetailsResponse ¶ added in v0.74.9
type GetRunDetailsResponse struct {
Input []byte `protobuf:"bytes,1,opt,name=input,proto3" json:"input,omitempty"` // the input payload for the workflow run
Status RunStatus `protobuf:"varint,2,opt,name=status,proto3,enum=v1.RunStatus" json:"status,omitempty"` // the status of the workflow run
TaskRuns map[string]*TaskRunDetail `` // map of task run external ids to their details
/* 173-byte string literal not displayed */
Done bool `protobuf:"varint,4,opt,name=done,proto3" json:"done,omitempty"` // indicates if the workflow run is done
AdditionalMetadata []byte `protobuf:"bytes,5,opt,name=additional_metadata,json=additionalMetadata,proto3" json:"additional_metadata,omitempty"` // (optional) additional metadata for the workflow run
IsEvicted bool `protobuf:"varint,6,opt,name=is_evicted,json=isEvicted,proto3" json:"is_evicted,omitempty"` // whether any task in this run has been evicted
// contains filtered or unexported fields
}
func (*GetRunDetailsResponse) Descriptor
deprecated
added in
v0.74.9
func (*GetRunDetailsResponse) Descriptor() ([]byte, []int)
Deprecated: Use GetRunDetailsResponse.ProtoReflect.Descriptor instead.
func (*GetRunDetailsResponse) GetAdditionalMetadata ¶ added in v0.74.13
func (x *GetRunDetailsResponse) GetAdditionalMetadata() []byte
func (*GetRunDetailsResponse) GetDone ¶ added in v0.74.9
func (x *GetRunDetailsResponse) GetDone() bool
func (*GetRunDetailsResponse) GetInput ¶ added in v0.74.9
func (x *GetRunDetailsResponse) GetInput() []byte
func (*GetRunDetailsResponse) GetIsEvicted ¶ added in v0.80.0
func (x *GetRunDetailsResponse) GetIsEvicted() bool
func (*GetRunDetailsResponse) GetStatus ¶ added in v0.74.9
func (x *GetRunDetailsResponse) GetStatus() RunStatus
func (*GetRunDetailsResponse) GetTaskRuns ¶ added in v0.74.9
func (x *GetRunDetailsResponse) GetTaskRuns() map[string]*TaskRunDetail
func (*GetRunDetailsResponse) ProtoMessage ¶ added in v0.74.9
func (*GetRunDetailsResponse) ProtoMessage()
func (*GetRunDetailsResponse) ProtoReflect ¶ added in v0.74.9
func (x *GetRunDetailsResponse) ProtoReflect() protoreflect.Message
func (*GetRunDetailsResponse) Reset ¶ added in v0.74.9
func (x *GetRunDetailsResponse) Reset()
func (*GetRunDetailsResponse) String ¶ added in v0.74.9
func (x *GetRunDetailsResponse) String() string
type GetStreamTopicMetadataRequest ¶ added in v0.110.13
type GetStreamTopicMetadataRequest struct {
Namespace string `protobuf:"bytes,1,opt,name=namespace,proto3" json:"namespace,omitempty"`
Topic string `protobuf:"bytes,2,opt,name=topic,proto3" json:"topic,omitempty"`
// contains filtered or unexported fields
}
func (*GetStreamTopicMetadataRequest) Descriptor
deprecated
added in
v0.110.13
func (*GetStreamTopicMetadataRequest) Descriptor() ([]byte, []int)
Deprecated: Use GetStreamTopicMetadataRequest.ProtoReflect.Descriptor instead.
func (*GetStreamTopicMetadataRequest) GetNamespace ¶ added in v0.110.13
func (x *GetStreamTopicMetadataRequest) GetNamespace() string
func (*GetStreamTopicMetadataRequest) GetTopic ¶ added in v0.110.13
func (x *GetStreamTopicMetadataRequest) GetTopic() string
func (*GetStreamTopicMetadataRequest) ProtoMessage ¶ added in v0.110.13
func (*GetStreamTopicMetadataRequest) ProtoMessage()
func (*GetStreamTopicMetadataRequest) ProtoReflect ¶ added in v0.110.13
func (x *GetStreamTopicMetadataRequest) ProtoReflect() protoreflect.Message
func (*GetStreamTopicMetadataRequest) Reset ¶ added in v0.110.13
func (x *GetStreamTopicMetadataRequest) Reset()
func (*GetStreamTopicMetadataRequest) String ¶ added in v0.110.13
func (x *GetStreamTopicMetadataRequest) String() string
type IdempotencyCollisionError ¶ added in v0.95.0
type IdempotencyCollisionError struct {
ExistingRunExternalId string `` // the external ID of the existing workflow run that caused the collision
/* 128-byte string literal not displayed */
CollidingRunExternalId string `` // the external ID of the workflow run that caused the collision
/* 131-byte string literal not displayed */
// contains filtered or unexported fields
}
func (*IdempotencyCollisionError) Descriptor
deprecated
added in
v0.95.0
func (*IdempotencyCollisionError) Descriptor() ([]byte, []int)
Deprecated: Use IdempotencyCollisionError.ProtoReflect.Descriptor instead.
func (*IdempotencyCollisionError) GetCollidingRunExternalId ¶ added in v0.95.0
func (x *IdempotencyCollisionError) GetCollidingRunExternalId() string
func (*IdempotencyCollisionError) GetExistingRunExternalId ¶ added in v0.95.0
func (x *IdempotencyCollisionError) GetExistingRunExternalId() string
func (*IdempotencyCollisionError) ProtoMessage ¶ added in v0.95.0
func (*IdempotencyCollisionError) ProtoMessage()
func (*IdempotencyCollisionError) ProtoReflect ¶ added in v0.95.0
func (x *IdempotencyCollisionError) ProtoReflect() protoreflect.Message
func (*IdempotencyCollisionError) Reset ¶ added in v0.95.0
func (x *IdempotencyCollisionError) Reset()
func (*IdempotencyCollisionError) String ¶ added in v0.95.0
func (x *IdempotencyCollisionError) String() string
type IdempotencyConfig ¶ added in v0.95.0
type IdempotencyConfig struct {
Expression string `protobuf:"bytes,1,opt,name=expression,proto3" json:"expression,omitempty"` // a CEL expression for determining the idempotency key for workflow runs
TtlMs int64 `protobuf:"varint,2,opt,name=ttl_ms,json=ttlMs,proto3" json:"ttl_ms,omitempty"` // time-to-live for idempotency keys in milliseconds. if the method is `STATUS`, this is a "fallback" - the longest the key can live before it's evicted
Method *IdempotencyMethod `protobuf:"varint,3,opt,name=method,proto3,enum=v1.IdempotencyMethod,oneof" json:"method,omitempty"` // the method to use for idempotency, defaults to TTL
// contains filtered or unexported fields
}
func (*IdempotencyConfig) Descriptor
deprecated
added in
v0.95.0
func (*IdempotencyConfig) Descriptor() ([]byte, []int)
Deprecated: Use IdempotencyConfig.ProtoReflect.Descriptor instead.
func (*IdempotencyConfig) GetExpression ¶ added in v0.95.0
func (x *IdempotencyConfig) GetExpression() string
func (*IdempotencyConfig) GetMethod ¶ added in v0.96.0
func (x *IdempotencyConfig) GetMethod() IdempotencyMethod
func (*IdempotencyConfig) GetTtlMs ¶ added in v0.95.0
func (x *IdempotencyConfig) GetTtlMs() int64
func (*IdempotencyConfig) ProtoMessage ¶ added in v0.95.0
func (*IdempotencyConfig) ProtoMessage()
func (*IdempotencyConfig) ProtoReflect ¶ added in v0.95.0
func (x *IdempotencyConfig) ProtoReflect() protoreflect.Message
func (*IdempotencyConfig) Reset ¶ added in v0.95.0
func (x *IdempotencyConfig) Reset()
func (*IdempotencyConfig) String ¶ added in v0.95.0
func (x *IdempotencyConfig) String() string
type IdempotencyMethod ¶ added in v0.96.0
type IdempotencyMethod int32
const ( IdempotencyMethod_TTL IdempotencyMethod = 0 IdempotencyMethod_STATUS IdempotencyMethod = 1 )
func (IdempotencyMethod) Descriptor ¶ added in v0.96.0
func (IdempotencyMethod) Descriptor() protoreflect.EnumDescriptor
func (IdempotencyMethod) Enum ¶ added in v0.96.0
func (x IdempotencyMethod) Enum() *IdempotencyMethod
func (IdempotencyMethod) EnumDescriptor
deprecated
added in
v0.96.0
func (IdempotencyMethod) EnumDescriptor() ([]byte, []int)
Deprecated: Use IdempotencyMethod.Descriptor instead.
func (IdempotencyMethod) Number ¶ added in v0.96.0
func (x IdempotencyMethod) Number() protoreflect.EnumNumber
func (IdempotencyMethod) String ¶ added in v0.96.0
func (x IdempotencyMethod) String() string
func (IdempotencyMethod) Type ¶ added in v0.96.0
func (IdempotencyMethod) Type() protoreflect.EnumType
type ListenForDurableEventRequest ¶
type ListenForDurableEventRequest struct {
TaskId string `protobuf:"bytes,1,opt,name=task_id,json=taskId,proto3" json:"task_id,omitempty"` // single listener per worker
SignalKey string `protobuf:"bytes,2,opt,name=signal_key,json=signalKey,proto3" json:"signal_key,omitempty"` // the match id for the listener
// contains filtered or unexported fields
}
func (*ListenForDurableEventRequest) Descriptor
deprecated
func (*ListenForDurableEventRequest) Descriptor() ([]byte, []int)
Deprecated: Use ListenForDurableEventRequest.ProtoReflect.Descriptor instead.
func (*ListenForDurableEventRequest) GetSignalKey ¶
func (x *ListenForDurableEventRequest) GetSignalKey() string
func (*ListenForDurableEventRequest) GetTaskId ¶
func (x *ListenForDurableEventRequest) GetTaskId() string
func (*ListenForDurableEventRequest) ProtoMessage ¶
func (*ListenForDurableEventRequest) ProtoMessage()
func (*ListenForDurableEventRequest) ProtoReflect ¶
func (x *ListenForDurableEventRequest) ProtoReflect() protoreflect.Message
func (*ListenForDurableEventRequest) Reset ¶
func (x *ListenForDurableEventRequest) Reset()
func (*ListenForDurableEventRequest) String ¶
func (x *ListenForDurableEventRequest) String() string
type OperatorActionsAck ¶ added in v0.109.0
type OperatorActionsAck struct {
Sequence uint64 `protobuf:"varint,1,opt,name=sequence,proto3" json:"sequence,omitempty"`
// contains filtered or unexported fields
}
OperatorActionsAck confirms that the delta with this sequence, and every delta with a lower sequence on the stream, has been committed to the worker's action set.
func (*OperatorActionsAck) Descriptor
deprecated
added in
v0.109.0
func (*OperatorActionsAck) Descriptor() ([]byte, []int)
Deprecated: Use OperatorActionsAck.ProtoReflect.Descriptor instead.
func (*OperatorActionsAck) GetSequence ¶ added in v0.109.0
func (x *OperatorActionsAck) GetSequence() uint64
func (*OperatorActionsAck) ProtoMessage ¶ added in v0.109.0
func (*OperatorActionsAck) ProtoMessage()
func (*OperatorActionsAck) ProtoReflect ¶ added in v0.109.0
func (x *OperatorActionsAck) ProtoReflect() protoreflect.Message
func (*OperatorActionsAck) Reset ¶ added in v0.109.0
func (x *OperatorActionsAck) Reset()
func (*OperatorActionsAck) String ¶ added in v0.109.0
func (x *OperatorActionsAck) String() string
type OperatorActionsDelta ¶ added in v0.109.0
type OperatorActionsDelta struct {
// Action ids the worker can now run.
Add []string `protobuf:"bytes,1,rep,name=add,proto3" json:"add,omitempty"`
// Action ids the worker no longer runs.
Remove []string `protobuf:"bytes,2,rep,name=remove,proto3" json:"remove,omitempty"`
// Sequence number of this delta on the stream: positive and strictly increasing. The
// server acknowledges the delta with an OperatorActionsAck carrying the same value once
// the change is committed. 0 asks for no ack.
Sequence uint64 `protobuf:"varint,3,opt,name=sequence,proto3" json:"sequence,omitempty"`
// contains filtered or unexported fields
}
OperatorActionsDelta changes the worker's action set incrementally. Adds are applied before removes. Adding an action the worker already has, or removing one it does not have, is a no-op. At most 1000 ids per message across both lists.
func (*OperatorActionsDelta) Descriptor
deprecated
added in
v0.109.0
func (*OperatorActionsDelta) Descriptor() ([]byte, []int)
Deprecated: Use OperatorActionsDelta.ProtoReflect.Descriptor instead.
func (*OperatorActionsDelta) GetAdd ¶ added in v0.109.0
func (x *OperatorActionsDelta) GetAdd() []string
func (*OperatorActionsDelta) GetRemove ¶ added in v0.109.0
func (x *OperatorActionsDelta) GetRemove() []string
func (*OperatorActionsDelta) GetSequence ¶ added in v0.109.0
func (x *OperatorActionsDelta) GetSequence() uint64
func (*OperatorActionsDelta) ProtoMessage ¶ added in v0.109.0
func (*OperatorActionsDelta) ProtoMessage()
func (*OperatorActionsDelta) ProtoReflect ¶ added in v0.109.0
func (x *OperatorActionsDelta) ProtoReflect() protoreflect.Message
func (*OperatorActionsDelta) Reset ¶ added in v0.109.0
func (x *OperatorActionsDelta) Reset()
func (*OperatorActionsDelta) String ¶ added in v0.109.0
func (x *OperatorActionsDelta) String() string
type OperatorHeartbeat ¶ added in v0.109.0
type OperatorHeartbeat struct {
HeartbeatAt *timestamppb.Timestamp `protobuf:"bytes,1,opt,name=heartbeat_at,json=heartbeatAt,proto3" json:"heartbeat_at,omitempty"`
// contains filtered or unexported fields
}
OperatorHeartbeat keeps the worker alive. Clients send one every 4 seconds on the Listen stream.
SDK workers heartbeat out of stream (Dispatcher.Heartbeat), from a thread of their own, because the runtimes the SDKs target can block their event loop while a task runs and would starve an in-stream heartbeat, marking a busy worker dead. Operators are written against the Go SDK, whose runtime is not subject to event loop blocking that way, so the heartbeat rides the stream it keeps alive: one connection, one liveness, and a stream that is gone takes its heartbeat with it.
func (*OperatorHeartbeat) Descriptor
deprecated
added in
v0.109.0
func (*OperatorHeartbeat) Descriptor() ([]byte, []int)
Deprecated: Use OperatorHeartbeat.ProtoReflect.Descriptor instead.
func (*OperatorHeartbeat) GetHeartbeatAt ¶ added in v0.109.0
func (x *OperatorHeartbeat) GetHeartbeatAt() *timestamppb.Timestamp
func (*OperatorHeartbeat) ProtoMessage ¶ added in v0.109.0
func (*OperatorHeartbeat) ProtoMessage()
func (*OperatorHeartbeat) ProtoReflect ¶ added in v0.109.0
func (x *OperatorHeartbeat) ProtoReflect() protoreflect.Message
func (*OperatorHeartbeat) Reset ¶ added in v0.109.0
func (x *OperatorHeartbeat) Reset()
func (*OperatorHeartbeat) String ¶ added in v0.109.0
func (x *OperatorHeartbeat) String() string
type OperatorListenRequest ¶ added in v0.109.0
type OperatorListenRequest struct {
// Types that are assignable to Message:
//
// *OperatorListenRequest_Start
// *OperatorListenRequest_Heartbeat
// *OperatorListenRequest_Actions
// *OperatorListenRequest_Pause
Message isOperatorListenRequest_Message `protobuf_oneof:"message"`
// contains filtered or unexported fields
}
OperatorListenRequest is a client-to-server message on the Listen stream. The first message must be start; every later message is a heartbeat, an actions delta or a pause.
func (*OperatorListenRequest) Descriptor
deprecated
added in
v0.109.0
func (*OperatorListenRequest) Descriptor() ([]byte, []int)
Deprecated: Use OperatorListenRequest.ProtoReflect.Descriptor instead.
func (*OperatorListenRequest) GetActions ¶ added in v0.109.0
func (x *OperatorListenRequest) GetActions() *OperatorActionsDelta
func (*OperatorListenRequest) GetHeartbeat ¶ added in v0.109.0
func (x *OperatorListenRequest) GetHeartbeat() *OperatorHeartbeat
func (*OperatorListenRequest) GetMessage ¶ added in v0.109.0
func (m *OperatorListenRequest) GetMessage() isOperatorListenRequest_Message
func (*OperatorListenRequest) GetPause ¶ added in v0.109.0
func (x *OperatorListenRequest) GetPause() *OperatorPause
func (*OperatorListenRequest) GetStart ¶ added in v0.109.0
func (x *OperatorListenRequest) GetStart() *OperatorListenStart
func (*OperatorListenRequest) ProtoMessage ¶ added in v0.109.0
func (*OperatorListenRequest) ProtoMessage()
func (*OperatorListenRequest) ProtoReflect ¶ added in v0.109.0
func (x *OperatorListenRequest) ProtoReflect() protoreflect.Message
func (*OperatorListenRequest) Reset ¶ added in v0.109.0
func (x *OperatorListenRequest) Reset()
func (*OperatorListenRequest) String ¶ added in v0.109.0
func (x *OperatorListenRequest) String() string
type OperatorListenRequest_Actions ¶ added in v0.109.0
type OperatorListenRequest_Actions struct {
Actions *OperatorActionsDelta `protobuf:"bytes,3,opt,name=actions,proto3,oneof"`
}
type OperatorListenRequest_Heartbeat ¶ added in v0.109.0
type OperatorListenRequest_Heartbeat struct {
Heartbeat *OperatorHeartbeat `protobuf:"bytes,2,opt,name=heartbeat,proto3,oneof"`
}
type OperatorListenRequest_Pause ¶ added in v0.109.0
type OperatorListenRequest_Pause struct {
Pause *OperatorPause `protobuf:"bytes,4,opt,name=pause,proto3,oneof"`
}
type OperatorListenRequest_Start ¶ added in v0.109.0
type OperatorListenRequest_Start struct {
Start *OperatorListenStart `protobuf:"bytes,1,opt,name=start,proto3,oneof"`
}
type OperatorListenResponse ¶ added in v0.109.0
type OperatorListenResponse struct {
// Types that are assignable to Message:
//
// *OperatorListenResponse_Action
// *OperatorListenResponse_Ack
// *OperatorListenResponse_PauseAck
Message isOperatorListenResponse_Message `protobuf_oneof:"message"`
// contains filtered or unexported fields
}
OperatorListenResponse is a server-to-client message on the Listen stream.
func (*OperatorListenResponse) Descriptor
deprecated
added in
v0.109.0
func (*OperatorListenResponse) Descriptor() ([]byte, []int)
Deprecated: Use OperatorListenResponse.ProtoReflect.Descriptor instead.
func (*OperatorListenResponse) GetAck ¶ added in v0.109.0
func (x *OperatorListenResponse) GetAck() *OperatorActionsAck
func (*OperatorListenResponse) GetAction ¶ added in v0.109.0
func (x *OperatorListenResponse) GetAction() *contracts.AssignedAction
func (*OperatorListenResponse) GetMessage ¶ added in v0.109.0
func (m *OperatorListenResponse) GetMessage() isOperatorListenResponse_Message
func (*OperatorListenResponse) GetPauseAck ¶ added in v0.109.0
func (x *OperatorListenResponse) GetPauseAck() *OperatorPauseAck
func (*OperatorListenResponse) ProtoMessage ¶ added in v0.109.0
func (*OperatorListenResponse) ProtoMessage()
func (*OperatorListenResponse) ProtoReflect ¶ added in v0.109.0
func (x *OperatorListenResponse) ProtoReflect() protoreflect.Message
func (*OperatorListenResponse) Reset ¶ added in v0.109.0
func (x *OperatorListenResponse) Reset()
func (*OperatorListenResponse) String ¶ added in v0.109.0
func (x *OperatorListenResponse) String() string
type OperatorListenResponse_Ack ¶ added in v0.109.0
type OperatorListenResponse_Ack struct {
// Acknowledgement of a committed actions delta.
Ack *OperatorActionsAck `protobuf:"bytes,2,opt,name=ack,proto3,oneof"`
}
type OperatorListenResponse_Action ¶ added in v0.109.0
type OperatorListenResponse_Action struct {
// An action assigned to the worker by the dispatcher.
Action *contracts.AssignedAction `protobuf:"bytes,1,opt,name=action,proto3,oneof"`
}
type OperatorListenResponse_PauseAck ¶ added in v0.109.0
type OperatorListenResponse_PauseAck struct {
// Acknowledgement of a committed pause.
PauseAck *OperatorPauseAck `protobuf:"bytes,3,opt,name=pause_ack,json=pauseAck,proto3,oneof"`
}
type OperatorListenStart ¶ added in v0.109.0
type OperatorListenStart struct {
// Worker id from OperatorRegisterResponse. The worker must belong to the calling operator.
WorkerId string `protobuf:"bytes,1,opt,name=worker_id,json=workerId,proto3" json:"worker_id,omitempty"`
// contains filtered or unexported fields
}
OperatorListenStart names the worker this stream activates.
func (*OperatorListenStart) Descriptor
deprecated
added in
v0.109.0
func (*OperatorListenStart) Descriptor() ([]byte, []int)
Deprecated: Use OperatorListenStart.ProtoReflect.Descriptor instead.
func (*OperatorListenStart) GetWorkerId ¶ added in v0.109.0
func (x *OperatorListenStart) GetWorkerId() string
func (*OperatorListenStart) ProtoMessage ¶ added in v0.109.0
func (*OperatorListenStart) ProtoMessage()
func (*OperatorListenStart) ProtoReflect ¶ added in v0.109.0
func (x *OperatorListenStart) ProtoReflect() protoreflect.Message
func (*OperatorListenStart) Reset ¶ added in v0.109.0
func (x *OperatorListenStart) Reset()
func (*OperatorListenStart) String ¶ added in v0.109.0
func (x *OperatorListenStart) String() string
type OperatorPause ¶ added in v0.109.0
type OperatorPause struct {
// True to stop the scheduler assigning to the worker, false to let it be assigned to again.
IsPaused bool `protobuf:"varint,1,opt,name=is_paused,json=isPaused,proto3" json:"is_paused,omitempty"`
// contains filtered or unexported fields
}
OperatorPause stops the scheduler assigning to the stream's worker, or lets it be assigned to again. The server answers with an OperatorPauseAck once the change is committed; after the ack of a pause nothing is delivered on the stream until the pause is lifted.
func (*OperatorPause) Descriptor
deprecated
added in
v0.109.0
func (*OperatorPause) Descriptor() ([]byte, []int)
Deprecated: Use OperatorPause.ProtoReflect.Descriptor instead.
func (*OperatorPause) GetIsPaused ¶ added in v0.109.0
func (x *OperatorPause) GetIsPaused() bool
func (*OperatorPause) ProtoMessage ¶ added in v0.109.0
func (*OperatorPause) ProtoMessage()
func (*OperatorPause) ProtoReflect ¶ added in v0.109.0
func (x *OperatorPause) ProtoReflect() protoreflect.Message
func (*OperatorPause) Reset ¶ added in v0.109.0
func (x *OperatorPause) Reset()
func (*OperatorPause) String ¶ added in v0.109.0
func (x *OperatorPause) String() string
type OperatorPauseAck ¶ added in v0.109.0
type OperatorPauseAck struct {
IsPaused bool `protobuf:"varint,1,opt,name=is_paused,json=isPaused,proto3" json:"is_paused,omitempty"`
// contains filtered or unexported fields
}
OperatorPauseAck confirms that the pause state is committed on the worker. For is_paused = true it also confirms that no action will be delivered on the stream from here on.
func (*OperatorPauseAck) Descriptor
deprecated
added in
v0.109.0
func (*OperatorPauseAck) Descriptor() ([]byte, []int)
Deprecated: Use OperatorPauseAck.ProtoReflect.Descriptor instead.
func (*OperatorPauseAck) GetIsPaused ¶ added in v0.109.0
func (x *OperatorPauseAck) GetIsPaused() bool
func (*OperatorPauseAck) ProtoMessage ¶ added in v0.109.0
func (*OperatorPauseAck) ProtoMessage()
func (*OperatorPauseAck) ProtoReflect ¶ added in v0.109.0
func (x *OperatorPauseAck) ProtoReflect() protoreflect.Message
func (*OperatorPauseAck) Reset ¶ added in v0.109.0
func (x *OperatorPauseAck) Reset()
func (*OperatorPauseAck) String ¶ added in v0.109.0
func (x *OperatorPauseAck) String() string
type OperatorRegisterRequest ¶ added in v0.109.0
type OperatorRegisterRequest struct {
// Operator name, unique per (tenant, kind = GRPC). The operator row is upserted by this name.
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
// Slot config for the worker (slot_type -> max units). Defaults to {"default": 100} when
// empty.
SlotConfig map[string]int32 `` /* 180-byte string literal not displayed */
// Worker labels used for affinity-based assignment.
Labels map[string]*contracts.WorkerLabels `` /* 153-byte string literal not displayed */
// Runtime information about the operator process, shown in the dashboard.
RuntimeInfo *contracts.RuntimeInfo `protobuf:"bytes,4,opt,name=runtime_info,json=runtimeInfo,proto3,oneof" json:"runtime_info,omitempty"`
// Worker id from a previous OperatorRegisterResponse on this operator. When set and the
// worker still exists for this operator, the response resumes that worker instead of
// creating a new one.
WorkerId *string `protobuf:"bytes,5,opt,name=worker_id,json=workerId,proto3,oneof" json:"worker_id,omitempty"`
// contains filtered or unexported fields
}
OperatorRegisterRequest identifies the operator and describes the worker for this connection.
func (*OperatorRegisterRequest) Descriptor
deprecated
added in
v0.109.0
func (*OperatorRegisterRequest) Descriptor() ([]byte, []int)
Deprecated: Use OperatorRegisterRequest.ProtoReflect.Descriptor instead.
func (*OperatorRegisterRequest) GetLabels ¶ added in v0.109.0
func (x *OperatorRegisterRequest) GetLabels() map[string]*contracts.WorkerLabels
func (*OperatorRegisterRequest) GetName ¶ added in v0.109.0
func (x *OperatorRegisterRequest) GetName() string
func (*OperatorRegisterRequest) GetRuntimeInfo ¶ added in v0.109.0
func (x *OperatorRegisterRequest) GetRuntimeInfo() *contracts.RuntimeInfo
func (*OperatorRegisterRequest) GetSlotConfig ¶ added in v0.109.0
func (x *OperatorRegisterRequest) GetSlotConfig() map[string]int32
func (*OperatorRegisterRequest) GetWorkerId ¶ added in v0.109.0
func (x *OperatorRegisterRequest) GetWorkerId() string
func (*OperatorRegisterRequest) ProtoMessage ¶ added in v0.109.0
func (*OperatorRegisterRequest) ProtoMessage()
func (*OperatorRegisterRequest) ProtoReflect ¶ added in v0.109.0
func (x *OperatorRegisterRequest) ProtoReflect() protoreflect.Message
func (*OperatorRegisterRequest) Reset ¶ added in v0.109.0
func (x *OperatorRegisterRequest) Reset()
func (*OperatorRegisterRequest) String ¶ added in v0.109.0
func (x *OperatorRegisterRequest) String() string
type OperatorRegisterResponse ¶ added in v0.109.0
type OperatorRegisterResponse struct {
TenantId string `protobuf:"bytes,1,opt,name=tenant_id,json=tenantId,proto3" json:"tenant_id,omitempty"`
// Operator id. Send it as the hatchet-operator-id gRPC metadata value on every later RPC.
OperatorId string `protobuf:"bytes,2,opt,name=operator_id,json=operatorId,proto3" json:"operator_id,omitempty"`
// Worker id backing this connection. Use it in the Listen start message, for
// SendStepActionEvent and DurableTask, and to resume the worker on reconnect.
WorkerId string `protobuf:"bytes,3,opt,name=worker_id,json=workerId,proto3" json:"worker_id,omitempty"`
// True when worker_id is the worker named in the request, false when a new worker was
// created (no worker_id was sent, or the worker no longer exists for this operator). A
// client that keeps its own view of the action set must replay it when resumed is false.
Resumed bool `protobuf:"varint,4,opt,name=resumed,proto3" json:"resumed,omitempty"`
// contains filtered or unexported fields
}
func (*OperatorRegisterResponse) Descriptor
deprecated
added in
v0.109.0
func (*OperatorRegisterResponse) Descriptor() ([]byte, []int)
Deprecated: Use OperatorRegisterResponse.ProtoReflect.Descriptor instead.
func (*OperatorRegisterResponse) GetOperatorId ¶ added in v0.109.0
func (x *OperatorRegisterResponse) GetOperatorId() string
func (*OperatorRegisterResponse) GetResumed ¶ added in v0.109.0
func (x *OperatorRegisterResponse) GetResumed() bool
func (*OperatorRegisterResponse) GetTenantId ¶ added in v0.109.0
func (x *OperatorRegisterResponse) GetTenantId() string
func (*OperatorRegisterResponse) GetWorkerId ¶ added in v0.109.0
func (x *OperatorRegisterResponse) GetWorkerId() string
func (*OperatorRegisterResponse) ProtoMessage ¶ added in v0.109.0
func (*OperatorRegisterResponse) ProtoMessage()
func (*OperatorRegisterResponse) ProtoReflect ¶ added in v0.109.0
func (x *OperatorRegisterResponse) ProtoReflect() protoreflect.Message
func (*OperatorRegisterResponse) Reset ¶ added in v0.109.0
func (x *OperatorRegisterResponse) Reset()
func (*OperatorRegisterResponse) String ¶ added in v0.109.0
func (x *OperatorRegisterResponse) String() string
type OperatorServiceClient ¶ added in v0.109.0
type OperatorServiceClient interface {
// Register upserts the operator and creates or resumes the worker for this connection. It
// does not register workflows or actions: workflows are put through AdminService.PutWorkflow
// on the same connection, and actions are added on the Listen stream.
Register(ctx context.Context, in *OperatorRegisterRequest, opts ...grpc.CallOption) (*OperatorRegisterResponse, error)
// Listen activates a registered worker for the lifetime of the stream, streams assigned
// actions to it, and accepts heartbeats and action set deltas from it.
//
// Protocol:
// 1. The first client message MUST be start, naming the worker returned by Register. Any
// other first message is rejected with InvalidArgument, and a second start on the same
// stream is also rejected. The server waits at most 30 seconds for it.
// 2. Every server message is an OperatorListenResponse: either an assigned action from the
// dispatcher fan-out or an ack for an actions delta.
// 3. The client sends a heartbeat every 4 seconds on this stream. A worker whose heartbeat
// goes stale is treated as inactive by the scheduler. See OperatorHeartbeat for why the
// heartbeat is in-stream rather than a separate RPC like an SDK worker's.
// 4. The client sends actions deltas whenever the set of actions it can run changes. Deltas
// are applied incrementally to the worker's action set and never replace it; a removed
// action stops being assigned to the worker within about a second. Each delta carries at
// most 1000 ids (adds plus removes); larger deltas are rejected with InvalidArgument.
// 5. Delta acknowledgement. Each delta carries a sequence number that is positive and
// strictly increasing on the stream. The server applies deltas in stream order and
// answers each one with an OperatorActionsAck carrying its sequence once the change is
// committed. An ack for sequence N therefore also confirms every lower sequence. A delta
// whose sequence is 0 is applied but never acknowledged. The server ends the stream
// instead of acknowledging when a delta cannot be applied or the ack cannot be sent.
// 6. Recovery. The client keeps every sent delta until its ack arrives. After a reconnect
// it first resends the unacknowledged deltas in order; if the reconnect did not resume
// the worker (OperatorRegisterResponse.resumed is false) it instead resends its whole
// desired action set, since a new worker starts with no actions. Applying a delta twice
// is harmless: adding an action the worker has and removing one it lacks are no-ops, so
// a delta that was committed but not acknowledged converges on replay.
// 7. When the stream ends (client close, network failure, or engine shutdown) the worker is
// deactivated. Reconnect by calling Register with worker_id set, then Listen again; the
// resumed worker keeps its action set.
// 8. Pause. The stream is stateful: a pause message stops the scheduler assigning to the
// worker, and from the moment the server answers with an OperatorPauseAck no action is
// delivered on the stream; one that the scheduler had already assigned is returned to
// the queue instead. An operator that pauses before draining therefore knows the
// actions it holds are the last it will get. pause with is_paused = false lets the worker
// be assigned to again and is acknowledged the same way. The pause belongs to the
// stream: Register clears it when it resumes the worker, so a client that reconnects
// while paused sends the pause again on the new stream, before its deltas.
//
// Requires the hatchet-operator-id metadata.
Listen(ctx context.Context, opts ...grpc.CallOption) (OperatorService_ListenClient, error)
// SendStepActionEvent reports task progress (started, completed, failed) for an action that
// was delivered on a Listen stream. Same semantics as Dispatcher.SendStepActionEvent.
// Requires the hatchet-operator-id metadata; the event's worker_id must belong to the
// operator.
SendStepActionEvent(ctx context.Context, in *contracts.StepActionEvent, opts ...grpc.CallOption) (*contracts.ActionEventResponse, error)
// DurableTask is the durable task event stream, identical to V1Dispatcher.DurableTask. The
// first client message registers the worker id returned by Register; a worker that belongs
// to another operator is rejected with PermissionDenied. Requires the hatchet-operator-id
// metadata.
DurableTask(ctx context.Context, opts ...grpc.CallOption) (OperatorService_DurableTaskClient, error)
}
OperatorServiceClient is the client API for OperatorService service.
For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
func NewOperatorServiceClient ¶ added in v0.109.0
func NewOperatorServiceClient(cc grpc.ClientConnInterface) OperatorServiceClient
type OperatorServiceServer ¶ added in v0.109.0
type OperatorServiceServer interface {
// Register upserts the operator and creates or resumes the worker for this connection. It
// does not register workflows or actions: workflows are put through AdminService.PutWorkflow
// on the same connection, and actions are added on the Listen stream.
Register(context.Context, *OperatorRegisterRequest) (*OperatorRegisterResponse, error)
// Listen activates a registered worker for the lifetime of the stream, streams assigned
// actions to it, and accepts heartbeats and action set deltas from it.
//
// Protocol:
// 1. The first client message MUST be start, naming the worker returned by Register. Any
// other first message is rejected with InvalidArgument, and a second start on the same
// stream is also rejected. The server waits at most 30 seconds for it.
// 2. Every server message is an OperatorListenResponse: either an assigned action from the
// dispatcher fan-out or an ack for an actions delta.
// 3. The client sends a heartbeat every 4 seconds on this stream. A worker whose heartbeat
// goes stale is treated as inactive by the scheduler. See OperatorHeartbeat for why the
// heartbeat is in-stream rather than a separate RPC like an SDK worker's.
// 4. The client sends actions deltas whenever the set of actions it can run changes. Deltas
// are applied incrementally to the worker's action set and never replace it; a removed
// action stops being assigned to the worker within about a second. Each delta carries at
// most 1000 ids (adds plus removes); larger deltas are rejected with InvalidArgument.
// 5. Delta acknowledgement. Each delta carries a sequence number that is positive and
// strictly increasing on the stream. The server applies deltas in stream order and
// answers each one with an OperatorActionsAck carrying its sequence once the change is
// committed. An ack for sequence N therefore also confirms every lower sequence. A delta
// whose sequence is 0 is applied but never acknowledged. The server ends the stream
// instead of acknowledging when a delta cannot be applied or the ack cannot be sent.
// 6. Recovery. The client keeps every sent delta until its ack arrives. After a reconnect
// it first resends the unacknowledged deltas in order; if the reconnect did not resume
// the worker (OperatorRegisterResponse.resumed is false) it instead resends its whole
// desired action set, since a new worker starts with no actions. Applying a delta twice
// is harmless: adding an action the worker has and removing one it lacks are no-ops, so
// a delta that was committed but not acknowledged converges on replay.
// 7. When the stream ends (client close, network failure, or engine shutdown) the worker is
// deactivated. Reconnect by calling Register with worker_id set, then Listen again; the
// resumed worker keeps its action set.
// 8. Pause. The stream is stateful: a pause message stops the scheduler assigning to the
// worker, and from the moment the server answers with an OperatorPauseAck no action is
// delivered on the stream; one that the scheduler had already assigned is returned to
// the queue instead. An operator that pauses before draining therefore knows the
// actions it holds are the last it will get. pause with is_paused = false lets the worker
// be assigned to again and is acknowledged the same way. The pause belongs to the
// stream: Register clears it when it resumes the worker, so a client that reconnects
// while paused sends the pause again on the new stream, before its deltas.
//
// Requires the hatchet-operator-id metadata.
Listen(OperatorService_ListenServer) error
// SendStepActionEvent reports task progress (started, completed, failed) for an action that
// was delivered on a Listen stream. Same semantics as Dispatcher.SendStepActionEvent.
// Requires the hatchet-operator-id metadata; the event's worker_id must belong to the
// operator.
SendStepActionEvent(context.Context, *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)
// DurableTask is the durable task event stream, identical to V1Dispatcher.DurableTask. The
// first client message registers the worker id returned by Register; a worker that belongs
// to another operator is rejected with PermissionDenied. Requires the hatchet-operator-id
// metadata.
DurableTask(OperatorService_DurableTaskServer) error
// contains filtered or unexported methods
}
OperatorServiceServer is the server API for OperatorService service. All implementations must embed UnimplementedOperatorServiceServer for forward compatibility
type OperatorService_DurableTaskClient ¶ added in v0.109.0
type OperatorService_DurableTaskClient interface {
Send(*DurableTaskRequest) error
Recv() (*DurableTaskResponse, error)
grpc.ClientStream
}
type OperatorService_DurableTaskServer ¶ added in v0.109.0
type OperatorService_DurableTaskServer interface {
Send(*DurableTaskResponse) error
Recv() (*DurableTaskRequest, error)
grpc.ServerStream
}
type OperatorService_ListenClient ¶ added in v0.109.0
type OperatorService_ListenClient interface {
Send(*OperatorListenRequest) error
Recv() (*OperatorListenResponse, error)
grpc.ClientStream
}
type OperatorService_ListenServer ¶ added in v0.109.0
type OperatorService_ListenServer interface {
Send(*OperatorListenResponse) error
Recv() (*OperatorListenRequest, error)
grpc.ServerStream
}
type ParentOverrideMatchCondition ¶
type ParentOverrideMatchCondition struct {
Base *BaseMatchCondition `protobuf:"bytes,1,opt,name=base,proto3" json:"base,omitempty"`
ParentReadableId string `protobuf:"bytes,2,opt,name=parent_readable_id,json=parentReadableId,proto3" json:"parent_readable_id,omitempty"`
// contains filtered or unexported fields
}
func (*ParentOverrideMatchCondition) Descriptor
deprecated
func (*ParentOverrideMatchCondition) Descriptor() ([]byte, []int)
Deprecated: Use ParentOverrideMatchCondition.ProtoReflect.Descriptor instead.
func (*ParentOverrideMatchCondition) GetBase ¶
func (x *ParentOverrideMatchCondition) GetBase() *BaseMatchCondition
func (*ParentOverrideMatchCondition) GetParentReadableId ¶
func (x *ParentOverrideMatchCondition) GetParentReadableId() string
func (*ParentOverrideMatchCondition) ProtoMessage ¶
func (*ParentOverrideMatchCondition) ProtoMessage()
func (*ParentOverrideMatchCondition) ProtoReflect ¶
func (x *ParentOverrideMatchCondition) ProtoReflect() protoreflect.Message
func (*ParentOverrideMatchCondition) Reset ¶
func (x *ParentOverrideMatchCondition) Reset()
func (*ParentOverrideMatchCondition) String ¶
func (x *ParentOverrideMatchCondition) String() string
type PublishStreamMessageRequest ¶ added in v0.110.0
type PublishStreamMessageRequest struct {
Namespace string `protobuf:"bytes,1,opt,name=namespace,proto3" json:"namespace,omitempty"`
Topic string `protobuf:"bytes,2,opt,name=topic,proto3" json:"topic,omitempty"`
Payload []byte `protobuf:"bytes,3,opt,name=payload,proto3" json:"payload,omitempty"`
ProducerId string `protobuf:"bytes,4,opt,name=producer_id,json=producerId,proto3" json:"producer_id,omitempty"`
ProducerSeq int64 `protobuf:"varint,5,opt,name=producer_seq,json=producerSeq,proto3" json:"producer_seq,omitempty"`
// instead of payload: a ref returned by the stream payload upload endpoint,
// for payloads too large for one gRPC message
PayloadRef string `protobuf:"bytes,6,opt,name=payload_ref,json=payloadRef,proto3" json:"payload_ref,omitempty"`
// contains filtered or unexported fields
}
func (*PublishStreamMessageRequest) Descriptor
deprecated
added in
v0.110.0
func (*PublishStreamMessageRequest) Descriptor() ([]byte, []int)
Deprecated: Use PublishStreamMessageRequest.ProtoReflect.Descriptor instead.
func (*PublishStreamMessageRequest) GetNamespace ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) GetNamespace() string
func (*PublishStreamMessageRequest) GetPayload ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) GetPayload() []byte
func (*PublishStreamMessageRequest) GetPayloadRef ¶ added in v0.110.13
func (x *PublishStreamMessageRequest) GetPayloadRef() string
func (*PublishStreamMessageRequest) GetProducerId ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) GetProducerId() string
func (*PublishStreamMessageRequest) GetProducerSeq ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) GetProducerSeq() int64
func (*PublishStreamMessageRequest) GetTopic ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) GetTopic() string
func (*PublishStreamMessageRequest) ProtoMessage ¶ added in v0.110.0
func (*PublishStreamMessageRequest) ProtoMessage()
func (*PublishStreamMessageRequest) ProtoReflect ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) ProtoReflect() protoreflect.Message
func (*PublishStreamMessageRequest) Reset ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) Reset()
func (*PublishStreamMessageRequest) String ¶ added in v0.110.0
func (x *PublishStreamMessageRequest) String() string
type PublishStreamMessageResponse ¶ added in v0.110.0
type PublishStreamMessageResponse struct {
// contains filtered or unexported fields
}
func (*PublishStreamMessageResponse) Descriptor
deprecated
added in
v0.110.0
func (*PublishStreamMessageResponse) Descriptor() ([]byte, []int)
Deprecated: Use PublishStreamMessageResponse.ProtoReflect.Descriptor instead.
func (*PublishStreamMessageResponse) ProtoMessage ¶ added in v0.110.0
func (*PublishStreamMessageResponse) ProtoMessage()
func (*PublishStreamMessageResponse) ProtoReflect ¶ added in v0.110.0
func (x *PublishStreamMessageResponse) ProtoReflect() protoreflect.Message
func (*PublishStreamMessageResponse) Reset ¶ added in v0.110.0
func (x *PublishStreamMessageResponse) Reset()
func (*PublishStreamMessageResponse) String ¶ added in v0.110.0
func (x *PublishStreamMessageResponse) String() string
type RateLimitDuration ¶
type RateLimitDuration int32
const ( RateLimitDuration_SECOND RateLimitDuration = 0 RateLimitDuration_MINUTE RateLimitDuration = 1 RateLimitDuration_HOUR RateLimitDuration = 2 RateLimitDuration_DAY RateLimitDuration = 3 RateLimitDuration_WEEK RateLimitDuration = 4 RateLimitDuration_MONTH RateLimitDuration = 5 RateLimitDuration_YEAR RateLimitDuration = 6 )
func (RateLimitDuration) Descriptor ¶
func (RateLimitDuration) Descriptor() protoreflect.EnumDescriptor
func (RateLimitDuration) Enum ¶
func (x RateLimitDuration) Enum() *RateLimitDuration
func (RateLimitDuration) EnumDescriptor
deprecated
func (RateLimitDuration) EnumDescriptor() ([]byte, []int)
Deprecated: Use RateLimitDuration.Descriptor instead.
func (RateLimitDuration) Number ¶
func (x RateLimitDuration) Number() protoreflect.EnumNumber
func (RateLimitDuration) String ¶
func (x RateLimitDuration) String() string
func (RateLimitDuration) Type ¶
func (RateLimitDuration) Type() protoreflect.EnumType
type RegisterDurableEventRequest ¶
type RegisterDurableEventRequest struct {
TaskId string `protobuf:"bytes,1,opt,name=task_id,json=taskId,proto3" json:"task_id,omitempty"` // external uuid for the task run
SignalKey string `protobuf:"bytes,2,opt,name=signal_key,json=signalKey,proto3" json:"signal_key,omitempty"` // the signal key for the event
Conditions *DurableEventListenerConditions `protobuf:"bytes,3,opt,name=conditions,proto3" json:"conditions,omitempty"` // the task conditions for creating the task
// contains filtered or unexported fields
}
func (*RegisterDurableEventRequest) Descriptor
deprecated
func (*RegisterDurableEventRequest) Descriptor() ([]byte, []int)
Deprecated: Use RegisterDurableEventRequest.ProtoReflect.Descriptor instead.
func (*RegisterDurableEventRequest) GetConditions ¶
func (x *RegisterDurableEventRequest) GetConditions() *DurableEventListenerConditions
func (*RegisterDurableEventRequest) GetSignalKey ¶
func (x *RegisterDurableEventRequest) GetSignalKey() string
func (*RegisterDurableEventRequest) GetTaskId ¶
func (x *RegisterDurableEventRequest) GetTaskId() string
func (*RegisterDurableEventRequest) ProtoMessage ¶
func (*RegisterDurableEventRequest) ProtoMessage()
func (*RegisterDurableEventRequest) ProtoReflect ¶
func (x *RegisterDurableEventRequest) ProtoReflect() protoreflect.Message
func (*RegisterDurableEventRequest) Reset ¶
func (x *RegisterDurableEventRequest) Reset()
func (*RegisterDurableEventRequest) String ¶
func (x *RegisterDurableEventRequest) String() string
type RegisterDurableEventResponse ¶
type RegisterDurableEventResponse struct {
// contains filtered or unexported fields
}
func (*RegisterDurableEventResponse) Descriptor
deprecated
func (*RegisterDurableEventResponse) Descriptor() ([]byte, []int)
Deprecated: Use RegisterDurableEventResponse.ProtoReflect.Descriptor instead.
func (*RegisterDurableEventResponse) ProtoMessage ¶
func (*RegisterDurableEventResponse) ProtoMessage()
func (*RegisterDurableEventResponse) ProtoReflect ¶
func (x *RegisterDurableEventResponse) ProtoReflect() protoreflect.Message
func (*RegisterDurableEventResponse) Reset ¶
func (x *RegisterDurableEventResponse) Reset()
func (*RegisterDurableEventResponse) String ¶
func (x *RegisterDurableEventResponse) String() string
type ReplayTasksRequest ¶
type ReplayTasksRequest struct {
ExternalIds []string `protobuf:"bytes,1,rep,name=external_ids,json=externalIds,proto3" json:"external_ids,omitempty"` // a list of external UUIDs
Filter *TasksFilter `protobuf:"bytes,2,opt,name=filter,proto3,oneof" json:"filter,omitempty"`
// contains filtered or unexported fields
}
func (*ReplayTasksRequest) Descriptor
deprecated
func (*ReplayTasksRequest) Descriptor() ([]byte, []int)
Deprecated: Use ReplayTasksRequest.ProtoReflect.Descriptor instead.
func (*ReplayTasksRequest) GetExternalIds ¶
func (x *ReplayTasksRequest) GetExternalIds() []string
func (*ReplayTasksRequest) GetFilter ¶
func (x *ReplayTasksRequest) GetFilter() *TasksFilter
func (*ReplayTasksRequest) ProtoMessage ¶
func (*ReplayTasksRequest) ProtoMessage()
func (*ReplayTasksRequest) ProtoReflect ¶
func (x *ReplayTasksRequest) ProtoReflect() protoreflect.Message
func (*ReplayTasksRequest) Reset ¶
func (x *ReplayTasksRequest) Reset()
func (*ReplayTasksRequest) String ¶
func (x *ReplayTasksRequest) String() string
type ReplayTasksResponse ¶
type ReplayTasksResponse struct {
ReplayedTasks []string `protobuf:"bytes,1,rep,name=replayed_tasks,json=replayedTasks,proto3" json:"replayed_tasks,omitempty"`
// contains filtered or unexported fields
}
func (*ReplayTasksResponse) Descriptor
deprecated
func (*ReplayTasksResponse) Descriptor() ([]byte, []int)
Deprecated: Use ReplayTasksResponse.ProtoReflect.Descriptor instead.
func (*ReplayTasksResponse) GetReplayedTasks ¶
func (x *ReplayTasksResponse) GetReplayedTasks() []string
func (*ReplayTasksResponse) ProtoMessage ¶
func (*ReplayTasksResponse) ProtoMessage()
func (*ReplayTasksResponse) ProtoReflect ¶
func (x *ReplayTasksResponse) ProtoReflect() protoreflect.Message
func (*ReplayTasksResponse) Reset ¶
func (x *ReplayTasksResponse) Reset()
func (*ReplayTasksResponse) String ¶
func (x *ReplayTasksResponse) String() string
type RunStatus ¶ added in v0.74.9
type RunStatus int32
func (RunStatus) Descriptor ¶ added in v0.74.9
func (RunStatus) Descriptor() protoreflect.EnumDescriptor
func (RunStatus) EnumDescriptor
deprecated
added in
v0.74.9
func (RunStatus) Number ¶ added in v0.74.9
func (x RunStatus) Number() protoreflect.EnumNumber
func (RunStatus) Type ¶ added in v0.74.9
func (RunStatus) Type() protoreflect.EnumType
type SleepMatchCondition ¶
type SleepMatchCondition struct {
Base *BaseMatchCondition `protobuf:"bytes,1,opt,name=base,proto3" json:"base,omitempty"`
SleepFor string `protobuf:"bytes,2,opt,name=sleep_for,json=sleepFor,proto3" json:"sleep_for,omitempty"` // a duration string indicating how long to sleep
// contains filtered or unexported fields
}
func (*SleepMatchCondition) Descriptor
deprecated
func (*SleepMatchCondition) Descriptor() ([]byte, []int)
Deprecated: Use SleepMatchCondition.ProtoReflect.Descriptor instead.
func (*SleepMatchCondition) GetBase ¶
func (x *SleepMatchCondition) GetBase() *BaseMatchCondition
func (*SleepMatchCondition) GetSleepFor ¶
func (x *SleepMatchCondition) GetSleepFor() string
func (*SleepMatchCondition) ProtoMessage ¶
func (*SleepMatchCondition) ProtoMessage()
func (*SleepMatchCondition) ProtoReflect ¶
func (x *SleepMatchCondition) ProtoReflect() protoreflect.Message
func (*SleepMatchCondition) Reset ¶
func (x *SleepMatchCondition) Reset()
func (*SleepMatchCondition) String ¶
func (x *SleepMatchCondition) String() string
type StickyStrategy ¶
type StickyStrategy int32
const ( StickyStrategy_SOFT StickyStrategy = 0 StickyStrategy_HARD StickyStrategy = 1 )
func (StickyStrategy) Descriptor ¶
func (StickyStrategy) Descriptor() protoreflect.EnumDescriptor
func (StickyStrategy) Enum ¶
func (x StickyStrategy) Enum() *StickyStrategy
func (StickyStrategy) EnumDescriptor
deprecated
func (StickyStrategy) EnumDescriptor() ([]byte, []int)
Deprecated: Use StickyStrategy.Descriptor instead.
func (StickyStrategy) Number ¶
func (x StickyStrategy) Number() protoreflect.EnumNumber
func (StickyStrategy) String ¶
func (x StickyStrategy) String() string
func (StickyStrategy) Type ¶
func (StickyStrategy) Type() protoreflect.EnumType
type StreamEntry ¶ added in v0.110.0
type StreamEntry struct {
Payload []byte `protobuf:"bytes,1,opt,name=payload,proto3" json:"payload,omitempty"`
Cursor string `protobuf:"bytes,2,opt,name=cursor,proto3" json:"cursor,omitempty"`
CreatedAt *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=created_at,json=createdAt,proto3" json:"created_at,omitempty"`
// instead of payload: fetch it from the stream payload endpoint by this ref
PayloadRef string `protobuf:"bytes,4,opt,name=payload_ref,json=payloadRef,proto3" json:"payload_ref,omitempty"`
// contains filtered or unexported fields
}
func (*StreamEntry) Descriptor
deprecated
added in
v0.110.0
func (*StreamEntry) Descriptor() ([]byte, []int)
Deprecated: Use StreamEntry.ProtoReflect.Descriptor instead.
func (*StreamEntry) GetCreatedAt ¶ added in v0.110.0
func (x *StreamEntry) GetCreatedAt() *timestamppb.Timestamp
func (*StreamEntry) GetCursor ¶ added in v0.110.0
func (x *StreamEntry) GetCursor() string
func (*StreamEntry) GetPayload ¶ added in v0.110.0
func (x *StreamEntry) GetPayload() []byte
func (*StreamEntry) GetPayloadRef ¶ added in v0.110.13
func (x *StreamEntry) GetPayloadRef() string
func (*StreamEntry) ProtoMessage ¶ added in v0.110.0
func (*StreamEntry) ProtoMessage()
func (*StreamEntry) ProtoReflect ¶ added in v0.110.0
func (x *StreamEntry) ProtoReflect() protoreflect.Message
func (*StreamEntry) Reset ¶ added in v0.110.0
func (x *StreamEntry) Reset()
func (*StreamEntry) String ¶ added in v0.110.0
func (x *StreamEntry) String() string
type StreamMessage ¶ added in v0.110.0
type StreamMessage struct {
Entries []*StreamEntry `protobuf:"bytes,1,rep,name=entries,proto3" json:"entries,omitempty"`
Hangup bool `protobuf:"varint,2,opt,name=hangup,proto3" json:"hangup,omitempty"`
Cursor string `protobuf:"bytes,3,opt,name=cursor,proto3" json:"cursor,omitempty"`
// contains filtered or unexported fields
}
func (*StreamMessage) Descriptor
deprecated
added in
v0.110.0
func (*StreamMessage) Descriptor() ([]byte, []int)
Deprecated: Use StreamMessage.ProtoReflect.Descriptor instead.
func (*StreamMessage) GetCursor ¶ added in v0.110.0
func (x *StreamMessage) GetCursor() string
func (*StreamMessage) GetEntries ¶ added in v0.110.0
func (x *StreamMessage) GetEntries() []*StreamEntry
func (*StreamMessage) GetHangup ¶ added in v0.110.0
func (x *StreamMessage) GetHangup() bool
func (*StreamMessage) ProtoMessage ¶ added in v0.110.0
func (*StreamMessage) ProtoMessage()
func (*StreamMessage) ProtoReflect ¶ added in v0.110.0
func (x *StreamMessage) ProtoReflect() protoreflect.Message
func (*StreamMessage) Reset ¶ added in v0.110.0
func (x *StreamMessage) Reset()
func (*StreamMessage) String ¶ added in v0.110.0
func (x *StreamMessage) String() string
type StreamTopicMetadata ¶ added in v0.110.13
type StreamTopicMetadata struct {
Namespace string `protobuf:"bytes,1,opt,name=namespace,proto3" json:"namespace,omitempty"`
Topic string `protobuf:"bytes,2,opt,name=topic,proto3" json:"topic,omitempty"`
TenantId string `protobuf:"bytes,3,opt,name=tenant_id,json=tenantId,proto3" json:"tenant_id,omitempty"`
// resumes a subscription after the newest retained message; unset when none is retained
LatestCursor *string `protobuf:"bytes,4,opt,name=latest_cursor,json=latestCursor,proto3,oneof" json:"latest_cursor,omitempty"`
MessageCount int64 `protobuf:"varint,5,opt,name=message_count,json=messageCount,proto3" json:"message_count,omitempty"`
// when the newest retained message was stored; unset when none is retained
LastPublishedAt *timestamppb.Timestamp `protobuf:"bytes,6,opt,name=last_published_at,json=lastPublishedAt,proto3" json:"last_published_at,omitempty"`
// contains filtered or unexported fields
}
Counts and the latest message cover only what's within the tenant's retention.
func (*StreamTopicMetadata) Descriptor
deprecated
added in
v0.110.13
func (*StreamTopicMetadata) Descriptor() ([]byte, []int)
Deprecated: Use StreamTopicMetadata.ProtoReflect.Descriptor instead.
func (*StreamTopicMetadata) GetLastPublishedAt ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetLastPublishedAt() *timestamppb.Timestamp
func (*StreamTopicMetadata) GetLatestCursor ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetLatestCursor() string
func (*StreamTopicMetadata) GetMessageCount ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetMessageCount() int64
func (*StreamTopicMetadata) GetNamespace ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetNamespace() string
func (*StreamTopicMetadata) GetTenantId ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetTenantId() string
func (*StreamTopicMetadata) GetTopic ¶ added in v0.110.13
func (x *StreamTopicMetadata) GetTopic() string
func (*StreamTopicMetadata) ProtoMessage ¶ added in v0.110.13
func (*StreamTopicMetadata) ProtoMessage()
func (*StreamTopicMetadata) ProtoReflect ¶ added in v0.110.13
func (x *StreamTopicMetadata) ProtoReflect() protoreflect.Message
func (*StreamTopicMetadata) Reset ¶ added in v0.110.13
func (x *StreamTopicMetadata) Reset()
func (*StreamTopicMetadata) String ¶ added in v0.110.13
func (x *StreamTopicMetadata) String() string
type SubscribeStreamRequest ¶ added in v0.110.0
type SubscribeStreamRequest struct {
Namespace string `protobuf:"bytes,1,opt,name=namespace,proto3" json:"namespace,omitempty"`
Topic string `protobuf:"bytes,2,opt,name=topic,proto3" json:"topic,omitempty"`
Cursor *string `protobuf:"bytes,3,opt,name=cursor,proto3,oneof" json:"cursor,omitempty"`
// contains filtered or unexported fields
}
func (*SubscribeStreamRequest) Descriptor
deprecated
added in
v0.110.0
func (*SubscribeStreamRequest) Descriptor() ([]byte, []int)
Deprecated: Use SubscribeStreamRequest.ProtoReflect.Descriptor instead.
func (*SubscribeStreamRequest) GetCursor ¶ added in v0.110.0
func (x *SubscribeStreamRequest) GetCursor() string
func (*SubscribeStreamRequest) GetNamespace ¶ added in v0.110.0
func (x *SubscribeStreamRequest) GetNamespace() string
func (*SubscribeStreamRequest) GetTopic ¶ added in v0.110.0
func (x *SubscribeStreamRequest) GetTopic() string
func (*SubscribeStreamRequest) ProtoMessage ¶ added in v0.110.0
func (*SubscribeStreamRequest) ProtoMessage()
func (*SubscribeStreamRequest) ProtoReflect ¶ added in v0.110.0
func (x *SubscribeStreamRequest) ProtoReflect() protoreflect.Message
func (*SubscribeStreamRequest) Reset ¶ added in v0.110.0
func (x *SubscribeStreamRequest) Reset()
func (*SubscribeStreamRequest) String ¶ added in v0.110.0
func (x *SubscribeStreamRequest) String() string
type TaskBatchConfig ¶ added in v0.98.0
type TaskBatchConfig struct {
BatchMaxSize int32 `protobuf:"varint,1,opt,name=batch_max_size,json=batchMaxSize,proto3" json:"batch_max_size,omitempty"` // (required) maximum items per batch
BatchMaxIntervalMs *int32 `` // (optional) time before batch flushes (milliseconds)
/* 126-byte string literal not displayed */
BatchGroupKey *string `protobuf:"bytes,3,opt,name=batch_group_key,json=batchGroupKey,proto3,oneof" json:"batch_group_key,omitempty"` // (optional) partition key for fairness (prevents mixing tenants)
BatchGroupMaxRuns *int32 `protobuf:"varint,4,opt,name=batch_group_max_runs,json=batchGroupMaxRuns,proto3,oneof" json:"batch_group_max_runs,omitempty"` // (optional) concurrent batches per group
BroadcastOutput *bool `protobuf:"varint,5,opt,name=broadcast_output,json=broadcastOutput,proto3,oneof" json:"broadcast_output,omitempty"` // (optional) when true, the handler returns one value broadcast to all callers; when false (default), the handler returns a dict keyed by step run id
// contains filtered or unexported fields
}
func (*TaskBatchConfig) Descriptor
deprecated
added in
v0.98.0
func (*TaskBatchConfig) Descriptor() ([]byte, []int)
Deprecated: Use TaskBatchConfig.ProtoReflect.Descriptor instead.
func (*TaskBatchConfig) GetBatchGroupKey ¶ added in v0.98.0
func (x *TaskBatchConfig) GetBatchGroupKey() string
func (*TaskBatchConfig) GetBatchGroupMaxRuns ¶ added in v0.98.0
func (x *TaskBatchConfig) GetBatchGroupMaxRuns() int32
func (*TaskBatchConfig) GetBatchMaxIntervalMs ¶ added in v0.98.0
func (x *TaskBatchConfig) GetBatchMaxIntervalMs() int32
func (*TaskBatchConfig) GetBatchMaxSize ¶ added in v0.98.0
func (x *TaskBatchConfig) GetBatchMaxSize() int32
func (*TaskBatchConfig) GetBroadcastOutput ¶ added in v0.98.0
func (x *TaskBatchConfig) GetBroadcastOutput() bool
func (*TaskBatchConfig) ProtoMessage ¶ added in v0.98.0
func (*TaskBatchConfig) ProtoMessage()
func (*TaskBatchConfig) ProtoReflect ¶ added in v0.98.0
func (x *TaskBatchConfig) ProtoReflect() protoreflect.Message
func (*TaskBatchConfig) Reset ¶ added in v0.98.0
func (x *TaskBatchConfig) Reset()
func (*TaskBatchConfig) String ¶ added in v0.98.0
func (x *TaskBatchConfig) String() string
type TaskConditions ¶
type TaskConditions struct {
ParentOverrideConditions []*ParentOverrideMatchCondition `` /* 135-byte string literal not displayed */
SleepConditions []*SleepMatchCondition `protobuf:"bytes,2,rep,name=sleep_conditions,json=sleepConditions,proto3" json:"sleep_conditions,omitempty"`
UserEventConditions []*UserEventMatchCondition `protobuf:"bytes,3,rep,name=user_event_conditions,json=userEventConditions,proto3" json:"user_event_conditions,omitempty"`
// contains filtered or unexported fields
}
func (*TaskConditions) Descriptor
deprecated
func (*TaskConditions) Descriptor() ([]byte, []int)
Deprecated: Use TaskConditions.ProtoReflect.Descriptor instead.
func (*TaskConditions) GetParentOverrideConditions ¶
func (x *TaskConditions) GetParentOverrideConditions() []*ParentOverrideMatchCondition
func (*TaskConditions) GetSleepConditions ¶
func (x *TaskConditions) GetSleepConditions() []*SleepMatchCondition
func (*TaskConditions) GetUserEventConditions ¶
func (x *TaskConditions) GetUserEventConditions() []*UserEventMatchCondition
func (*TaskConditions) ProtoMessage ¶
func (*TaskConditions) ProtoMessage()
func (*TaskConditions) ProtoReflect ¶
func (x *TaskConditions) ProtoReflect() protoreflect.Message
func (*TaskConditions) Reset ¶
func (x *TaskConditions) Reset()
func (*TaskConditions) String ¶
func (x *TaskConditions) String() string
type TaskRunDetail ¶ added in v0.74.9
type TaskRunDetail struct {
ExternalId string `protobuf:"bytes,1,opt,name=external_id,json=externalId,proto3" json:"external_id,omitempty"` // the external id (uuid) of the task run
Status RunStatus `protobuf:"varint,2,opt,name=status,proto3,enum=v1.RunStatus" json:"status,omitempty"` // the status of the task run
Error *string `protobuf:"bytes,3,opt,name=error,proto3,oneof" json:"error,omitempty"` // (optional) error message from the task run, if any
Output []byte `protobuf:"bytes,4,opt,name=output,proto3,oneof" json:"output,omitempty"` // (optional) the output payload for the task run
ReadableId string `protobuf:"bytes,5,opt,name=readable_id,json=readableId,proto3" json:"readable_id,omitempty"` // the readable id of the task
IsEvicted bool `protobuf:"varint,6,opt,name=is_evicted,json=isEvicted,proto3" json:"is_evicted,omitempty"` // whether the task has been evicted from a worker (status will be RUNNING)
// contains filtered or unexported fields
}
func (*TaskRunDetail) Descriptor
deprecated
added in
v0.74.9
func (*TaskRunDetail) Descriptor() ([]byte, []int)
Deprecated: Use TaskRunDetail.ProtoReflect.Descriptor instead.
func (*TaskRunDetail) GetError ¶ added in v0.74.9
func (x *TaskRunDetail) GetError() string
func (*TaskRunDetail) GetExternalId ¶ added in v0.74.9
func (x *TaskRunDetail) GetExternalId() string
func (*TaskRunDetail) GetIsEvicted ¶ added in v0.80.0
func (x *TaskRunDetail) GetIsEvicted() bool
func (*TaskRunDetail) GetOutput ¶ added in v0.74.9
func (x *TaskRunDetail) GetOutput() []byte
func (*TaskRunDetail) GetReadableId ¶ added in v0.74.9
func (x *TaskRunDetail) GetReadableId() string
func (*TaskRunDetail) GetStatus ¶ added in v0.74.9
func (x *TaskRunDetail) GetStatus() RunStatus
func (*TaskRunDetail) ProtoMessage ¶ added in v0.74.9
func (*TaskRunDetail) ProtoMessage()
func (*TaskRunDetail) ProtoReflect ¶ added in v0.74.9
func (x *TaskRunDetail) ProtoReflect() protoreflect.Message
func (*TaskRunDetail) Reset ¶ added in v0.74.9
func (x *TaskRunDetail) Reset()
func (*TaskRunDetail) String ¶ added in v0.74.9
func (x *TaskRunDetail) String() string
type TasksFilter ¶
type TasksFilter struct {
Statuses []string `protobuf:"bytes,1,rep,name=statuses,proto3" json:"statuses,omitempty"`
Since *timestamppb.Timestamp `protobuf:"bytes,2,opt,name=since,proto3" json:"since,omitempty"`
Until *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=until,proto3,oneof" json:"until,omitempty"`
WorkflowIds []string `protobuf:"bytes,4,rep,name=workflow_ids,json=workflowIds,proto3" json:"workflow_ids,omitempty"`
AdditionalMetadata []string `protobuf:"bytes,5,rep,name=additional_metadata,json=additionalMetadata,proto3" json:"additional_metadata,omitempty"`
// contains filtered or unexported fields
}
func (*TasksFilter) Descriptor
deprecated
func (*TasksFilter) Descriptor() ([]byte, []int)
Deprecated: Use TasksFilter.ProtoReflect.Descriptor instead.
func (*TasksFilter) GetAdditionalMetadata ¶
func (x *TasksFilter) GetAdditionalMetadata() []string
func (*TasksFilter) GetSince ¶
func (x *TasksFilter) GetSince() *timestamppb.Timestamp
func (*TasksFilter) GetStatuses ¶
func (x *TasksFilter) GetStatuses() []string
func (*TasksFilter) GetUntil ¶
func (x *TasksFilter) GetUntil() *timestamppb.Timestamp
func (*TasksFilter) GetWorkflowIds ¶
func (x *TasksFilter) GetWorkflowIds() []string
func (*TasksFilter) ProtoMessage ¶
func (*TasksFilter) ProtoMessage()
func (*TasksFilter) ProtoReflect ¶
func (x *TasksFilter) ProtoReflect() protoreflect.Message
func (*TasksFilter) Reset ¶
func (x *TasksFilter) Reset()
func (*TasksFilter) String ¶
func (x *TasksFilter) String() string
type TriggerWorkflowRequest ¶ added in v0.80.0
type TriggerWorkflowRequest struct {
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
// (optional) the input data for the workflow
Input string `protobuf:"bytes,2,opt,name=input,proto3" json:"input,omitempty"`
// (optional) the parent workflow run id
ParentId *string `protobuf:"bytes,3,opt,name=parent_id,json=parentId,proto3,oneof" json:"parent_id,omitempty"`
// (optional) the parent task external run id
ParentTaskRunExternalId *string `` /* 142-byte string literal not displayed */
// (optional) the index of the child workflow. if this is set, matches on the index or the
// child key will return an existing workflow run if the parent id, parent task run id, and
// child index/key match an existing workflow run.
ChildIndex *int32 `protobuf:"varint,5,opt,name=child_index,json=childIndex,proto3,oneof" json:"child_index,omitempty"`
// (optional) the key for the child. if this is set, matches on the index or the
// child key will return an existing workflow run if the parent id, parent task run id, and
// child index/key match an existing workflow run.
ChildKey *string `protobuf:"bytes,6,opt,name=child_key,json=childKey,proto3,oneof" json:"child_key,omitempty"`
// (optional) additional metadata for the workflow
AdditionalMetadata *string `protobuf:"bytes,7,opt,name=additional_metadata,json=additionalMetadata,proto3,oneof" json:"additional_metadata,omitempty"`
// (optional) desired worker id for the workflow run,
// requires the workflow definition to have a sticky strategy
DesiredWorkerId *string `protobuf:"bytes,8,opt,name=desired_worker_id,json=desiredWorkerId,proto3,oneof" json:"desired_worker_id,omitempty"`
// (optional) override for the priority of the workflow tasks, will set all tasks to this priority
Priority *int32 `protobuf:"varint,9,opt,name=priority,proto3,oneof" json:"priority,omitempty"`
// (optional) the desired worker labels for the workflow run, which will be used to determine which workers can pick up the workflow's tasks. if not set, defaults to an empty set of labels, which means any worker can pick up the tasks.
DesiredWorkerLabels map[string]*DesiredWorkerLabels `` /* 209-byte string literal not displayed */
// contains filtered or unexported fields
}
func (*TriggerWorkflowRequest) Descriptor
deprecated
added in
v0.80.0
func (*TriggerWorkflowRequest) Descriptor() ([]byte, []int)
Deprecated: Use TriggerWorkflowRequest.ProtoReflect.Descriptor instead.
func (*TriggerWorkflowRequest) GetAdditionalMetadata ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetAdditionalMetadata() string
func (*TriggerWorkflowRequest) GetChildIndex ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetChildIndex() int32
func (*TriggerWorkflowRequest) GetChildKey ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetChildKey() string
func (*TriggerWorkflowRequest) GetDesiredWorkerId ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetDesiredWorkerId() string
func (*TriggerWorkflowRequest) GetDesiredWorkerLabels ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetDesiredWorkerLabels() map[string]*DesiredWorkerLabels
func (*TriggerWorkflowRequest) GetInput ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetInput() string
func (*TriggerWorkflowRequest) GetName ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetName() string
func (*TriggerWorkflowRequest) GetParentId ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetParentId() string
func (*TriggerWorkflowRequest) GetParentTaskRunExternalId ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetParentTaskRunExternalId() string
func (*TriggerWorkflowRequest) GetPriority ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) GetPriority() int32
func (*TriggerWorkflowRequest) ProtoMessage ¶ added in v0.80.0
func (*TriggerWorkflowRequest) ProtoMessage()
func (*TriggerWorkflowRequest) ProtoReflect ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) ProtoReflect() protoreflect.Message
func (*TriggerWorkflowRequest) Reset ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) Reset()
func (*TriggerWorkflowRequest) String ¶ added in v0.80.0
func (x *TriggerWorkflowRequest) String() string
type TriggerWorkflowRunRequest ¶
type TriggerWorkflowRunRequest struct {
WorkflowName string `protobuf:"bytes,1,opt,name=workflow_name,json=workflowName,proto3" json:"workflow_name,omitempty"`
Input []byte `protobuf:"bytes,2,opt,name=input,proto3" json:"input,omitempty"`
AdditionalMetadata []byte `protobuf:"bytes,3,opt,name=additional_metadata,json=additionalMetadata,proto3" json:"additional_metadata,omitempty"`
Priority *int32 `protobuf:"varint,4,opt,name=priority,proto3,oneof" json:"priority,omitempty"`
DesiredWorkerLabels map[string]*DesiredWorkerLabels `` /* 208-byte string literal not displayed */
// contains filtered or unexported fields
}
func (*TriggerWorkflowRunRequest) Descriptor
deprecated
func (*TriggerWorkflowRunRequest) Descriptor() ([]byte, []int)
Deprecated: Use TriggerWorkflowRunRequest.ProtoReflect.Descriptor instead.
func (*TriggerWorkflowRunRequest) GetAdditionalMetadata ¶
func (x *TriggerWorkflowRunRequest) GetAdditionalMetadata() []byte
func (*TriggerWorkflowRunRequest) GetDesiredWorkerLabels ¶ added in v0.79.15
func (x *TriggerWorkflowRunRequest) GetDesiredWorkerLabels() map[string]*DesiredWorkerLabels
func (*TriggerWorkflowRunRequest) GetInput ¶
func (x *TriggerWorkflowRunRequest) GetInput() []byte
func (*TriggerWorkflowRunRequest) GetPriority ¶
func (x *TriggerWorkflowRunRequest) GetPriority() int32
func (*TriggerWorkflowRunRequest) GetWorkflowName ¶
func (x *TriggerWorkflowRunRequest) GetWorkflowName() string
func (*TriggerWorkflowRunRequest) ProtoMessage ¶
func (*TriggerWorkflowRunRequest) ProtoMessage()
func (*TriggerWorkflowRunRequest) ProtoReflect ¶
func (x *TriggerWorkflowRunRequest) ProtoReflect() protoreflect.Message
func (*TriggerWorkflowRunRequest) Reset ¶
func (x *TriggerWorkflowRunRequest) Reset()
func (*TriggerWorkflowRunRequest) String ¶
func (x *TriggerWorkflowRunRequest) String() string
type TriggerWorkflowRunResponse ¶
type TriggerWorkflowRunResponse struct {
ExternalId string `protobuf:"bytes,1,opt,name=external_id,json=externalId,proto3" json:"external_id,omitempty"`
// contains filtered or unexported fields
}
func (*TriggerWorkflowRunResponse) Descriptor
deprecated
func (*TriggerWorkflowRunResponse) Descriptor() ([]byte, []int)
Deprecated: Use TriggerWorkflowRunResponse.ProtoReflect.Descriptor instead.
func (*TriggerWorkflowRunResponse) GetExternalId ¶
func (x *TriggerWorkflowRunResponse) GetExternalId() string
func (*TriggerWorkflowRunResponse) ProtoMessage ¶
func (*TriggerWorkflowRunResponse) ProtoMessage()
func (*TriggerWorkflowRunResponse) ProtoReflect ¶
func (x *TriggerWorkflowRunResponse) ProtoReflect() protoreflect.Message
func (*TriggerWorkflowRunResponse) Reset ¶
func (x *TriggerWorkflowRunResponse) Reset()
func (*TriggerWorkflowRunResponse) String ¶
func (x *TriggerWorkflowRunResponse) String() string
type UnimplementedAdminServiceServer ¶
type UnimplementedAdminServiceServer struct {
}
UnimplementedAdminServiceServer must be embedded to have forward compatible implementations.
func (UnimplementedAdminServiceServer) BranchDurableTask ¶ added in v0.80.0
func (UnimplementedAdminServiceServer) BranchDurableTask(context.Context, *BranchDurableTaskRequest) (*BranchDurableTaskResponse, error)
func (UnimplementedAdminServiceServer) CancelTasks ¶
func (UnimplementedAdminServiceServer) CancelTasks(context.Context, *CancelTasksRequest) (*CancelTasksResponse, error)
func (UnimplementedAdminServiceServer) GetRunDetails ¶ added in v0.74.9
func (UnimplementedAdminServiceServer) GetRunDetails(context.Context, *GetRunDetailsRequest) (*GetRunDetailsResponse, error)
func (UnimplementedAdminServiceServer) PutWorkflow ¶
func (UnimplementedAdminServiceServer) PutWorkflow(context.Context, *CreateWorkflowVersionRequest) (*CreateWorkflowVersionResponse, error)
func (UnimplementedAdminServiceServer) ReplayTasks ¶
func (UnimplementedAdminServiceServer) ReplayTasks(context.Context, *ReplayTasksRequest) (*ReplayTasksResponse, error)
func (UnimplementedAdminServiceServer) TriggerWorkflowRun ¶
func (UnimplementedAdminServiceServer) TriggerWorkflowRun(context.Context, *TriggerWorkflowRunRequest) (*TriggerWorkflowRunResponse, error)
type UnimplementedOperatorServiceServer ¶ added in v0.109.0
type UnimplementedOperatorServiceServer struct {
}
UnimplementedOperatorServiceServer must be embedded to have forward compatible implementations.
func (UnimplementedOperatorServiceServer) DurableTask ¶ added in v0.109.0
func (UnimplementedOperatorServiceServer) DurableTask(OperatorService_DurableTaskServer) error
func (UnimplementedOperatorServiceServer) Listen ¶ added in v0.109.0
func (UnimplementedOperatorServiceServer) Listen(OperatorService_ListenServer) error
func (UnimplementedOperatorServiceServer) Register ¶ added in v0.109.0
func (UnimplementedOperatorServiceServer) Register(context.Context, *OperatorRegisterRequest) (*OperatorRegisterResponse, error)
func (UnimplementedOperatorServiceServer) SendStepActionEvent ¶ added in v0.109.0
func (UnimplementedOperatorServiceServer) SendStepActionEvent(context.Context, *contracts.StepActionEvent) (*contracts.ActionEventResponse, error)
type UnimplementedV1DispatcherServer ¶
type UnimplementedV1DispatcherServer struct {
}
UnimplementedV1DispatcherServer must be embedded to have forward compatible implementations.
func (UnimplementedV1DispatcherServer) DurableTask ¶ added in v0.80.0
func (UnimplementedV1DispatcherServer) DurableTask(V1Dispatcher_DurableTaskServer) error
func (UnimplementedV1DispatcherServer) ListenForDurableEvent ¶
func (UnimplementedV1DispatcherServer) ListenForDurableEvent(V1Dispatcher_ListenForDurableEventServer) error
func (UnimplementedV1DispatcherServer) RegisterDurableEvent ¶
func (UnimplementedV1DispatcherServer) RegisterDurableEvent(context.Context, *RegisterDurableEventRequest) (*RegisterDurableEventResponse, error)
type UnimplementedV1StreamsServer ¶ added in v0.110.0
type UnimplementedV1StreamsServer struct {
}
UnimplementedV1StreamsServer must be embedded to have forward compatible implementations.
func (UnimplementedV1StreamsServer) GetTopicMetadata ¶ added in v0.110.13
func (UnimplementedV1StreamsServer) GetTopicMetadata(context.Context, *GetStreamTopicMetadataRequest) (*StreamTopicMetadata, error)
func (UnimplementedV1StreamsServer) Publish ¶ added in v0.110.0
func (UnimplementedV1StreamsServer) Publish(context.Context, *PublishStreamMessageRequest) (*PublishStreamMessageResponse, error)
func (UnimplementedV1StreamsServer) Subscribe ¶ added in v0.110.0
func (UnimplementedV1StreamsServer) Subscribe(*SubscribeStreamRequest, V1Streams_SubscribeServer) error
type UnsafeAdminServiceServer ¶
type UnsafeAdminServiceServer interface {
// contains filtered or unexported methods
}
UnsafeAdminServiceServer may be embedded to opt out of forward compatibility for this service. Use of this interface is not recommended, as added methods to AdminServiceServer will result in compilation errors.
type UnsafeOperatorServiceServer ¶ added in v0.109.0
type UnsafeOperatorServiceServer interface {
// contains filtered or unexported methods
}
UnsafeOperatorServiceServer may be embedded to opt out of forward compatibility for this service. Use of this interface is not recommended, as added methods to OperatorServiceServer will result in compilation errors.
type UnsafeV1DispatcherServer ¶
type UnsafeV1DispatcherServer interface {
// contains filtered or unexported methods
}
UnsafeV1DispatcherServer may be embedded to opt out of forward compatibility for this service. Use of this interface is not recommended, as added methods to V1DispatcherServer will result in compilation errors.
type UnsafeV1StreamsServer ¶ added in v0.110.0
type UnsafeV1StreamsServer interface {
// contains filtered or unexported methods
}
UnsafeV1StreamsServer may be embedded to opt out of forward compatibility for this service. Use of this interface is not recommended, as added methods to V1StreamsServer will result in compilation errors.
type UserEventMatchCondition ¶
type UserEventMatchCondition struct {
Base *BaseMatchCondition `protobuf:"bytes,1,opt,name=base,proto3" json:"base,omitempty"`
UserEventKey string `protobuf:"bytes,2,opt,name=user_event_key,json=userEventKey,proto3" json:"user_event_key,omitempty"`
EventScope *string `protobuf:"bytes,3,opt,name=event_scope,json=eventScope,proto3,oneof" json:"event_scope,omitempty"` // an optional scope for the user event condition (similar to scopes on event filters)
ConsiderEventsSince *timestamppb.Timestamp `` /* 126-byte string literal not displayed */
// contains filtered or unexported fields
}
func (*UserEventMatchCondition) Descriptor
deprecated
func (*UserEventMatchCondition) Descriptor() ([]byte, []int)
Deprecated: Use UserEventMatchCondition.ProtoReflect.Descriptor instead.
func (*UserEventMatchCondition) GetBase ¶
func (x *UserEventMatchCondition) GetBase() *BaseMatchCondition
func (*UserEventMatchCondition) GetConsiderEventsSince ¶ added in v0.83.18
func (x *UserEventMatchCondition) GetConsiderEventsSince() *timestamppb.Timestamp
func (*UserEventMatchCondition) GetEventScope ¶ added in v0.83.18
func (x *UserEventMatchCondition) GetEventScope() string
func (*UserEventMatchCondition) GetUserEventKey ¶
func (x *UserEventMatchCondition) GetUserEventKey() string
func (*UserEventMatchCondition) ProtoMessage ¶
func (*UserEventMatchCondition) ProtoMessage()
func (*UserEventMatchCondition) ProtoReflect ¶
func (x *UserEventMatchCondition) ProtoReflect() protoreflect.Message
func (*UserEventMatchCondition) Reset ¶
func (x *UserEventMatchCondition) Reset()
func (*UserEventMatchCondition) String ¶
func (x *UserEventMatchCondition) String() string
type V1DispatcherClient ¶
type V1DispatcherClient interface {
DurableTask(ctx context.Context, opts ...grpc.CallOption) (V1Dispatcher_DurableTaskClient, error)
// NOTE: deprecated after DurableEventLog is implemented
RegisterDurableEvent(ctx context.Context, in *RegisterDurableEventRequest, opts ...grpc.CallOption) (*RegisterDurableEventResponse, error)
ListenForDurableEvent(ctx context.Context, opts ...grpc.CallOption) (V1Dispatcher_ListenForDurableEventClient, error)
}
V1DispatcherClient is the client API for V1Dispatcher service.
For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
func NewV1DispatcherClient ¶
func NewV1DispatcherClient(cc grpc.ClientConnInterface) V1DispatcherClient
type V1DispatcherServer ¶
type V1DispatcherServer interface {
DurableTask(V1Dispatcher_DurableTaskServer) error
// NOTE: deprecated after DurableEventLog is implemented
RegisterDurableEvent(context.Context, *RegisterDurableEventRequest) (*RegisterDurableEventResponse, error)
ListenForDurableEvent(V1Dispatcher_ListenForDurableEventServer) error
// contains filtered or unexported methods
}
V1DispatcherServer is the server API for V1Dispatcher service. All implementations must embed UnimplementedV1DispatcherServer for forward compatibility
type V1Dispatcher_DurableTaskClient ¶ added in v0.80.0
type V1Dispatcher_DurableTaskClient interface {
Send(*DurableTaskRequest) error
Recv() (*DurableTaskResponse, error)
grpc.ClientStream
}
type V1Dispatcher_DurableTaskServer ¶ added in v0.80.0
type V1Dispatcher_DurableTaskServer interface {
Send(*DurableTaskResponse) error
Recv() (*DurableTaskRequest, error)
grpc.ServerStream
}
type V1Dispatcher_ListenForDurableEventClient ¶
type V1Dispatcher_ListenForDurableEventClient interface {
Send(*ListenForDurableEventRequest) error
Recv() (*DurableEvent, error)
grpc.ClientStream
}
type V1Dispatcher_ListenForDurableEventServer ¶
type V1Dispatcher_ListenForDurableEventServer interface {
Send(*DurableEvent) error
Recv() (*ListenForDurableEventRequest, error)
grpc.ServerStream
}
type V1StreamsClient ¶ added in v0.110.0
type V1StreamsClient interface {
Publish(ctx context.Context, in *PublishStreamMessageRequest, opts ...grpc.CallOption) (*PublishStreamMessageResponse, error)
Subscribe(ctx context.Context, in *SubscribeStreamRequest, opts ...grpc.CallOption) (V1Streams_SubscribeClient, error)
GetTopicMetadata(ctx context.Context, in *GetStreamTopicMetadataRequest, opts ...grpc.CallOption) (*StreamTopicMetadata, error)
}
V1StreamsClient is the client API for V1Streams service.
For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
func NewV1StreamsClient ¶ added in v0.110.0
func NewV1StreamsClient(cc grpc.ClientConnInterface) V1StreamsClient
type V1StreamsServer ¶ added in v0.110.0
type V1StreamsServer interface {
Publish(context.Context, *PublishStreamMessageRequest) (*PublishStreamMessageResponse, error)
Subscribe(*SubscribeStreamRequest, V1Streams_SubscribeServer) error
GetTopicMetadata(context.Context, *GetStreamTopicMetadataRequest) (*StreamTopicMetadata, error)
// contains filtered or unexported methods
}
V1StreamsServer is the server API for V1Streams service. All implementations must embed UnimplementedV1StreamsServer for forward compatibility
type V1Streams_SubscribeClient ¶ added in v0.110.0
type V1Streams_SubscribeClient interface {
Recv() (*StreamMessage, error)
grpc.ClientStream
}
type V1Streams_SubscribeServer ¶ added in v0.110.0
type V1Streams_SubscribeServer interface {
Send(*StreamMessage) error
grpc.ServerStream
}
type WorkerLabelComparator ¶
type WorkerLabelComparator int32
const ( WorkerLabelComparator_EQUAL WorkerLabelComparator = 0 WorkerLabelComparator_NOT_EQUAL WorkerLabelComparator = 1 WorkerLabelComparator_GREATER_THAN WorkerLabelComparator = 2 WorkerLabelComparator_GREATER_THAN_OR_EQUAL WorkerLabelComparator = 3 WorkerLabelComparator_LESS_THAN WorkerLabelComparator = 4 WorkerLabelComparator_LESS_THAN_OR_EQUAL WorkerLabelComparator = 5 )
func (WorkerLabelComparator) Descriptor ¶
func (WorkerLabelComparator) Descriptor() protoreflect.EnumDescriptor
func (WorkerLabelComparator) Enum ¶
func (x WorkerLabelComparator) Enum() *WorkerLabelComparator
func (WorkerLabelComparator) EnumDescriptor
deprecated
func (WorkerLabelComparator) EnumDescriptor() ([]byte, []int)
Deprecated: Use WorkerLabelComparator.Descriptor instead.
func (WorkerLabelComparator) Number ¶
func (x WorkerLabelComparator) Number() protoreflect.EnumNumber
func (WorkerLabelComparator) String ¶
func (x WorkerLabelComparator) String() string
func (WorkerLabelComparator) Type ¶
func (WorkerLabelComparator) Type() protoreflect.EnumType