v1

package
v0.110.13 Latest Latest
Warning

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

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

Documentation

Index

Constants

This section is empty.

Variables

View Source
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.

View Source
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.

View Source
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.

View Source
var (
	StickyStrategy_name = map[int32]string{
		0: "SOFT",
		1: "HARD",
	}
	StickyStrategy_value = map[string]int32{
		"SOFT": 0,
		"HARD": 1,
	}
)

Enum value maps for StickyStrategy.

View Source
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.

View Source
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.

View Source
var (
	IdempotencyMethod_name = map[int32]string{
		0: "TTL",
		1: "STATUS",
	}
	IdempotencyMethod_value = map[string]int32{
		"TTL":    0,
		"STATUS": 1,
	}
)

Enum value maps for IdempotencyMethod.

View Source
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.

View Source
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)

View Source
var File_v1_dispatcher_proto protoreflect.FileDescriptor
View Source
var File_v1_operator_proto protoreflect.FileDescriptor
View Source
var File_v1_shared_condition_proto protoreflect.FileDescriptor
View Source
var File_v1_shared_trigger_proto protoreflect.FileDescriptor
View Source
var File_v1_streams_proto protoreflect.FileDescriptor
View Source
var File_v1_workflows_proto protoreflect.FileDescriptor
View Source
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)

View Source
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)

View Source
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
const (
	Action_CREATE Action = 0
	Action_QUEUE  Action = 1
	Action_CANCEL Action = 2
	Action_SKIP   Action = 3
)

func (Action) Descriptor

func (Action) Descriptor() protoreflect.EnumDescriptor

func (Action) Enum

func (x Action) Enum() *Action

func (Action) EnumDescriptor deprecated

func (Action) EnumDescriptor() ([]byte, []int)

Deprecated: Use Action.Descriptor instead.

func (Action) Number

func (x Action) Number() protoreflect.EnumNumber

func (Action) String

func (x Action) String() string

func (Action) Type

func (Action) Type() protoreflect.EnumType

type AdminServiceClient

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.

type AdminServiceServer

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 (*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 (*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 (*BulkTriggerIdempotencyCollisionError) Reset added in v0.95.0

func (*BulkTriggerIdempotencyCollisionError) String added in v0.95.0

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) Enum

func (ConcurrencyLimitStrategy) EnumDescriptor deprecated

func (ConcurrencyLimitStrategy) EnumDescriptor() ([]byte, []int)

Deprecated: Use ConcurrencyLimitStrategy.Descriptor instead.

func (ConcurrencyLimitStrategy) Number

func (ConcurrencyLimitStrategy) String

func (x ConcurrencyLimitStrategy) String() string

func (ConcurrencyLimitStrategy) Type

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 (*CreateWorkflowVersionRequest) GetTasks

func (*CreateWorkflowVersionRequest) GetVersion

func (x *CreateWorkflowVersionRequest) GetVersion() string

func (*CreateWorkflowVersionRequest) ProtoMessage

func (*CreateWorkflowVersionRequest) ProtoMessage()

func (*CreateWorkflowVersionRequest) ProtoReflect

func (*CreateWorkflowVersionRequest) Reset

func (x *CreateWorkflowVersionRequest) Reset()

func (*CreateWorkflowVersionRequest) 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 (*CreateWorkflowVersionResponse) GetWorkflowId

func (x *CreateWorkflowVersionResponse) GetWorkflowId() string

func (*CreateWorkflowVersionResponse) ProtoMessage

func (*CreateWorkflowVersionResponse) ProtoMessage()

func (*CreateWorkflowVersionResponse) ProtoReflect

func (*CreateWorkflowVersionResponse) Reset

func (x *CreateWorkflowVersionResponse) Reset()

func (*CreateWorkflowVersionResponse) 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 (*DurableEventListenerConditions) Reset

func (x *DurableEventListenerConditions) Reset()

func (*DurableEventListenerConditions) 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 (*DurableTaskAwaitedCompletedEntry) Reset added in v0.80.0

func (*DurableTaskAwaitedCompletedEntry) String added in v0.80.0

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 (*DurableTaskCompleteMemoRequest) ProtoMessage added in v0.80.0

func (*DurableTaskCompleteMemoRequest) ProtoMessage()

func (*DurableTaskCompleteMemoRequest) ProtoReflect added in v0.80.0

func (*DurableTaskCompleteMemoRequest) Reset added in v0.80.0

func (x *DurableTaskCompleteMemoRequest) Reset()

func (*DurableTaskCompleteMemoRequest) String added in v0.80.0

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 (*DurableTaskErrorResponse) GetRef added in v0.80.0

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) Enum added in v0.80.0

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 (DurableTaskErrorType) String added in v0.80.0

func (x DurableTaskErrorType) String() string

func (DurableTaskErrorType) Type added in v0.80.0

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 (*DurableTaskEventLogEntryCompletedResponse) GetIsFailure added in v0.89.3

func (*DurableTaskEventLogEntryCompletedResponse) GetPayload added in v0.80.0

func (*DurableTaskEventLogEntryCompletedResponse) GetRef added in v0.80.0

func (*DurableTaskEventLogEntryCompletedResponse) ProtoMessage added in v0.80.0

func (*DurableTaskEventLogEntryCompletedResponse) ProtoReflect added in v0.80.0

func (*DurableTaskEventLogEntryCompletedResponse) Reset added in v0.80.0

func (*DurableTaskEventLogEntryCompletedResponse) String added in v0.80.0

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 (*DurableTaskEventMemoAckResponse) ProtoMessage added in v0.80.0

func (*DurableTaskEventMemoAckResponse) ProtoMessage()

func (*DurableTaskEventMemoAckResponse) ProtoReflect added in v0.80.0

func (*DurableTaskEventMemoAckResponse) Reset added in v0.80.0

func (*DurableTaskEventMemoAckResponse) String added in v0.80.0

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 (*DurableTaskEventTriggerRunsAckResponse) ProtoMessage added in v0.80.0

func (*DurableTaskEventTriggerRunsAckResponse) ProtoReflect added in v0.80.0

func (*DurableTaskEventTriggerRunsAckResponse) Reset added in v0.80.0

func (*DurableTaskEventTriggerRunsAckResponse) String added in v0.80.0

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 (*DurableTaskEventWaitForAckResponse) ProtoMessage added in v0.80.0

func (*DurableTaskEventWaitForAckResponse) ProtoMessage()

func (*DurableTaskEventWaitForAckResponse) ProtoReflect added in v0.80.0

func (*DurableTaskEventWaitForAckResponse) Reset added in v0.80.0

func (*DurableTaskEventWaitForAckResponse) String added in v0.80.0

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 (*DurableTaskEvictInvocationRequest) ProtoMessage added in v0.80.0

func (*DurableTaskEvictInvocationRequest) ProtoMessage()

func (*DurableTaskEvictInvocationRequest) ProtoReflect added in v0.80.0

func (*DurableTaskEvictInvocationRequest) Reset added in v0.80.0

func (*DurableTaskEvictInvocationRequest) String added in v0.80.0

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 (*DurableTaskEvictionAckResponse) Reset added in v0.80.0

func (x *DurableTaskEvictionAckResponse) Reset()

func (*DurableTaskEvictionAckResponse) String added in v0.80.0

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 (*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 (*DurableTaskRequest) GetWaitFor added in v0.80.0

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 (*DurableTaskRequestRegisterWorker) Reset added in v0.80.0

func (*DurableTaskRequestRegisterWorker) String added in v0.80.0

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 (*DurableTaskResponse) GetError added in v0.80.0

func (*DurableTaskResponse) GetEvictionAck added in v0.80.0

func (*DurableTaskResponse) GetMemoAck added in v0.80.0

func (*DurableTaskResponse) GetMessage added in v0.80.0

func (m *DurableTaskResponse) GetMessage() isDurableTaskResponse_Message

func (*DurableTaskResponse) GetRegisterWorker added in v0.80.0

func (*DurableTaskResponse) GetServerEvict added in v0.80.0

func (*DurableTaskResponse) GetTriggerRunsAck added in v0.80.0

func (*DurableTaskResponse) GetWaitForAck added in v0.80.0

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 (*DurableTaskResponseRegisterWorker) Reset added in v0.80.0

func (*DurableTaskResponseRegisterWorker) String added in v0.80.0

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 (*DurableTaskServerEvictNotice) Reset added in v0.80.0

func (x *DurableTaskServerEvictNotice) Reset()

func (*DurableTaskServerEvictNotice) String added in v0.80.0

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 (*DurableTaskTriggerRunsRequest) ProtoMessage added in v0.80.0

func (*DurableTaskTriggerRunsRequest) ProtoMessage()

func (*DurableTaskTriggerRunsRequest) ProtoReflect added in v0.80.0

func (*DurableTaskTriggerRunsRequest) Reset added in v0.80.0

func (x *DurableTaskTriggerRunsRequest) Reset()

func (*DurableTaskTriggerRunsRequest) String added in v0.80.0

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 (*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 (*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 (*DurableTaskWorkerStatusRequest) Reset added in v0.80.0

func (x *DurableTaskWorkerStatusRequest) Reset()

func (*DurableTaskWorkerStatusRequest) String added in v0.80.0

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 (*GetStreamTopicMetadataRequest) Reset added in v0.110.13

func (x *GetStreamTopicMetadataRequest) Reset()

func (*GetStreamTopicMetadataRequest) String added in v0.110.13

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 (*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) Enum added in v0.96.0

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 (IdempotencyMethod) String added in v0.96.0

func (x IdempotencyMethod) String() string

func (IdempotencyMethod) Type added in v0.96.0

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 (*ListenForDurableEventRequest) Reset

func (x *ListenForDurableEventRequest) Reset()

func (*ListenForDurableEventRequest) 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 (*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 (*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 (*OperatorListenResponse) GetAction added in v0.109.0

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 (*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 (*ParentOverrideMatchCondition) GetParentReadableId

func (x *ParentOverrideMatchCondition) GetParentReadableId() string

func (*ParentOverrideMatchCondition) ProtoMessage

func (*ParentOverrideMatchCondition) ProtoMessage()

func (*ParentOverrideMatchCondition) ProtoReflect

func (*ParentOverrideMatchCondition) Reset

func (x *ParentOverrideMatchCondition) Reset()

func (*ParentOverrideMatchCondition) 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 (*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 (*PublishStreamMessageResponse) Reset added in v0.110.0

func (x *PublishStreamMessageResponse) Reset()

func (*PublishStreamMessageResponse) String added in v0.110.0

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) Enum

func (RateLimitDuration) EnumDescriptor deprecated

func (RateLimitDuration) EnumDescriptor() ([]byte, []int)

Deprecated: Use RateLimitDuration.Descriptor instead.

func (RateLimitDuration) Number

func (RateLimitDuration) String

func (x RateLimitDuration) String() string

func (RateLimitDuration) Type

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 (*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 (*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 (*RegisterDurableEventResponse) Reset

func (x *RegisterDurableEventResponse) Reset()

func (*RegisterDurableEventResponse) 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
const (
	RunStatus_QUEUED    RunStatus = 0
	RunStatus_RUNNING   RunStatus = 1
	RunStatus_COMPLETED RunStatus = 2
	RunStatus_FAILED    RunStatus = 3
	RunStatus_CANCELLED RunStatus = 4
	RunStatus_EVICTED   RunStatus = 5
)

func (RunStatus) Descriptor added in v0.74.9

func (RunStatus) Descriptor() protoreflect.EnumDescriptor

func (RunStatus) Enum added in v0.74.9

func (x RunStatus) Enum() *RunStatus

func (RunStatus) EnumDescriptor deprecated added in v0.74.9

func (RunStatus) EnumDescriptor() ([]byte, []int)

Deprecated: Use RunStatus.Descriptor instead.

func (RunStatus) Number added in v0.74.9

func (x RunStatus) Number() protoreflect.EnumNumber

func (RunStatus) String added in v0.74.9

func (x RunStatus) String() string

func (RunStatus) Type added in v0.74.9

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 (*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) Enum

func (x StickyStrategy) Enum() *StickyStrategy

func (StickyStrategy) EnumDescriptor deprecated

func (StickyStrategy) EnumDescriptor() ([]byte, []int)

Deprecated: Use StickyStrategy.Descriptor instead.

func (StickyStrategy) Number

func (StickyStrategy) String

func (x StickyStrategy) String() string

func (StickyStrategy) Type

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 (*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 (*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) CancelTasks

func (UnimplementedAdminServiceServer) GetRunDetails added in v0.74.9

func (UnimplementedAdminServiceServer) ReplayTasks

func (UnimplementedAdminServiceServer) TriggerWorkflowRun

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) Listen added in v0.109.0

func (UnimplementedOperatorServiceServer) Register added in v0.109.0

func (UnimplementedOperatorServiceServer) SendStepActionEvent added in v0.109.0

type UnimplementedV1DispatcherServer

type UnimplementedV1DispatcherServer struct {
}

UnimplementedV1DispatcherServer must be embedded to have forward compatible implementations.

func (UnimplementedV1DispatcherServer) DurableTask added in v0.80.0

func (UnimplementedV1DispatcherServer) ListenForDurableEvent

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) Publish added in v0.110.0

func (UnimplementedV1StreamsServer) Subscribe added in v0.110.0

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 (*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.

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

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) Enum

func (WorkerLabelComparator) EnumDescriptor deprecated

func (WorkerLabelComparator) EnumDescriptor() ([]byte, []int)

Deprecated: Use WorkerLabelComparator.Descriptor instead.

func (WorkerLabelComparator) Number

func (WorkerLabelComparator) String

func (x WorkerLabelComparator) String() string

func (WorkerLabelComparator) Type

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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