go.temporal.io/server/service/matching/handler.go
726 LOC · 253 covered · 473 uncovered · 39 ranges · 73 concepts · 15 introducers · 17 tests
File neighbourhood
The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.
Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file
In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.
Graph controls are ready.
Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.
package matching
import (
"context"
"sync"
"time"
enumspb "go.temporal.io/api/enums/v1"
taskqueuepb "go.temporal.io/api/taskqueue/v1"
workerpb "go.temporal.io/api/worker/v1"
"go.temporal.io/api/workflowservice/v1"
"go.temporal.io/server/api/matchingservice/v1"
"go.temporal.io/server/common"
"go.temporal.io/server/common/cluster"
"go.temporal.io/server/common/headers"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/membership"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/namespace"
"go.temporal.io/server/common/persistence"
"go.temporal.io/server/common/persistence/serialization"
"go.temporal.io/server/common/persistence/visibility/manager"
"go.temporal.io/server/common/primitives"
"go.temporal.io/server/common/resource"
"go.temporal.io/server/common/searchattribute"
"go.temporal.io/server/common/testing/testhooks"
"go.temporal.io/server/common/tqid"
"go.temporal.io/server/service/matching/hooks"
"go.temporal.io/server/service/matching/workers"
"go.temporal.io/server/service/worker/workerdeployment"
"go.uber.org/fx"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/reflect/protoreflect"
)
type (
// Handler - gRPC handler interface for matchingservice
Handler struct {
matchingservice.UnimplementedMatchingServiceServer
engine Engine
config *Config
metricsHandler metrics.Handler
logger log.Logger
startWG sync.WaitGroup
throttledLogger log.Logger
namespaceRegistry namespace.Registry
workersRegistry workers.Registry
}
HandlerParams struct {
fx.In
Config *Config
Logger log.Logger
ThrottledLogger log.Logger
TaskManager persistence.TaskManager
FairTaskManager persistence.FairTaskManager
HistoryClient resource.HistoryClient
MatchingRawClient resource.MatchingRawClient
WorkerDeploymentClient workerdeployment.Client
HostInfoProvider membership.HostInfoProvider
MatchingServiceResolver membership.ServiceResolver
MetricsHandler metrics.Handler
NamespaceRegistry namespace.Registry
ClusterMetadata cluster.Metadata
NamespaceReplicationQueue persistence.NamespaceReplicationQueue
VisibilityManager manager.VisibilityManager
NexusEndpointManager persistence.NexusEndpointManager
TestHooks testhooks.TestHooks
SearchAttributeProvider searchattribute.Provider
SearchAttributeMapperProvider searchattribute.MapperProvider
RateLimiter TaskDispatchRateLimiter `optional:"true"`
WorkersRegistry workers.Registry
Serializer serialization.Serializer
TaskHookFactories []hooks.TaskHookFactory `group:"TaskHookFactories"`
PartitionScalerFactory PartitionScalerFactory
}
)
const (
serviceName = "temporal.api.workflowservice.v1.MatchingService"
)
var (
_ matchingservice.MatchingServiceServer = (*Handler)(nil)
)
// NewHandler creates a gRPC handler for the matchingservice
func NewHandler(
params HandlerParams,
handler := &Handler{
config: params.Config,
metricsHandler: params.MetricsHandler,
logger: params.Logger,
throttledLogger: params.ThrottledLogger,
engine: NewEngine(
params.TaskManager,
params.FairTaskManager,
params.HistoryClient,
params.MatchingRawClient, // Use non retry client inside matching
params.WorkerDeploymentClient,
params.Config,
params.Logger,
params.ThrottledLogger,
params.MetricsHandler,
params.NamespaceRegistry,
params.HostInfoProvider,
params.MatchingServiceResolver,
params.ClusterMetadata,
params.NamespaceReplicationQueue,
params.VisibilityManager,
params.NexusEndpointManager,
params.TestHooks,
params.SearchAttributeProvider,
params.SearchAttributeMapperProvider,
params.RateLimiter,
params.Serializer,
params.TaskHookFactories,
params.PartitionScalerFactory,
),
namespaceRegistry: params.NamespaceRegistry,
workersRegistry: params.WorkersRegistry,
}
// prevent from serving requests before matching engine is started and ready
handler.startWG.Add(1)
return handler
}
// Start starts the handler
h.engine.Start()
h.startWG.Done()
}
// Stop stops the handler
h.engine.Stop()
}
func (h *Handler) opMetricsHandler(
namespaceID string,
taskQueue *taskqueuepb.TaskQueue,
taskQueueType enumspb.TaskQueueType,
operation string,
nsName := h.namespaceName(namespace.ID(namespaceID))
partition := tqid.UnsafePartitionFromProto(taskQueue, namespaceID, taskQueueType)
return metrics.GetPerTaskQueuePartitionIDScope(
h.metricsHandler.WithTags(metrics.OperationTag(operation)),
nsName.String(),
partition,
h.config.BreakdownMetricsByTaskQueue(nsName.String(), partition.TaskQueue().Name(), partition.TaskType()),
h.config.BreakdownMetricsByPartition(nsName.String(), partition.TaskQueue().Name(), partition.TaskType()),
)
}
// recordNexusTaskRequest emits the nexus_task_requests metric with namespace,
// operation, client_name, and is_internal tags.
func (h *Handler) recordNexusTaskRequest(ctx context.Context, namespaceID string, taskQueueKind enumspb.TaskQueueKind, operation string) {
tags.go ×3
nsName := h.namespaceName(namespace.ID(namespaceID))
clientName, _ := headers.GetClientNameAndVersion(ctx)
isInternal := primitives.IsInternalTaskQueueKind(taskQueueKind)
metrics.NexusTaskRequests.With(h.metricsHandler).Record(1,
metrics.NamespaceTag(nsName.String()),
metrics.OperationTag(operation),
metrics.ClientNameTag(clientName),
metrics.IsInternalTag(isInternal),
)
}
// AddActivityTask - adds an activity task.
func (h *Handler) AddActivityTask(
ctx context.Context,
request *matchingservice.AddActivityTaskRequest,
defer log.CapturePanic(h.logger, &retError)
startT := time.Now().UTC()
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_ACTIVITY,
metrics.MatchingAddActivityTaskScope,
)
if request.GetForwardInfo() != nil {
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
}
assignedBuildId, syncMatch, err := h.engine.AddActivityTask(ctx, request)
request_response.pb.go ×6
if syncMatch {
metrics.SyncMatchLatencyPerTaskQueue.With(opMetrics).Record(time.Since(startT))
}
return &matchingservice.AddActivityTaskResponse{AssignedBuildId: assignedBuildId}, err
}
// AddWorkflowTask - adds a workflow task.
func (h *Handler) AddWorkflowTask(
ctx context.Context,
request *matchingservice.AddWorkflowTaskRequest,
defer log.CapturePanic(h.logger, &retError)
startT := time.Now().UTC()
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_WORKFLOW,
metrics.MatchingAddWorkflowTaskScope,
)
if request.GetForwardInfo() != nil {
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
handler.go ×1
}
if syncMatch {
}
return &matchingservice.AddWorkflowTaskResponse{AssignedBuildId: assignedBuildId}, err
handler.go ×25
}
// PollActivityTaskQueue - long poll for an activity task.
func (h *Handler) PollActivityTaskQueue(
ctx context.Context,
request *matchingservice.PollActivityTaskQueueRequest,
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetPollRequest().GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_ACTIVITY,
metrics.MatchingPollActivityTaskQueueScope,
)
if request.GetForwardedSource() != "" {
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
request_response.pb.go ×6
}
ctx,
"PollActivityTaskQueue",
h.throttledLogger,
); err != nil {
return nil, err
}
}
// PollWorkflowTaskQueue - long poll for a workflow task.
func (h *Handler) PollWorkflowTaskQueue(
ctx context.Context,
request *matchingservice.PollWorkflowTaskQueueRequest,
) (_ *matchingservice.PollWorkflowTaskQueueResponseWithRawHistory, retError error) {
service_grpc.pb.go ×20
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetPollRequest().GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_WORKFLOW,
metrics.MatchingPollWorkflowTaskQueueScope,
)
if request.GetForwardedSource() != "" {
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
request_response.pb.go ×6
}
ctx,
"PollWorkflowTaskQueue",
h.throttledLogger,
); err != nil {
return nil, err
}
}
// QueryWorkflow queries a given workflow synchronously and return the query result.
func (h *Handler) QueryWorkflow(
ctx context.Context,
request *matchingservice.QueryWorkflowRequest,
) (_ *matchingservice.QueryWorkflowResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_WORKFLOW,
metrics.MatchingQueryWorkflowScope,
)
if request.GetForwardInfo() != nil {
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
}
return h.engine.QueryWorkflow(ctx, request)
}
// RespondQueryTaskCompleted responds a query task completed
func (h *Handler) RespondQueryTaskCompleted(
ctx context.Context,
request *matchingservice.RespondQueryTaskCompletedRequest,
) (_ *matchingservice.RespondQueryTaskCompletedResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_WORKFLOW,
metrics.MatchingRespondQueryTaskCompletedScope,
)
err := h.engine.RespondQueryTaskCompleted(ctx, request, opMetrics)
return &matchingservice.RespondQueryTaskCompletedResponse{}, err
}
// CancelOutstandingPoll is used to cancel outstanding pollers
func (h *Handler) CancelOutstandingPoll(ctx context.Context,
request *matchingservice.CancelOutstandingPollRequest) (_ *matchingservice.CancelOutstandingPollResponse, retError error) {
service_grpc.pb.go ×19
defer log.CapturePanic(h.logger, &retError)
err := h.engine.CancelOutstandingPoll(ctx, request)
return &matchingservice.CancelOutstandingPollResponse{}, err
}
// CancelOutstandingWorkerPolls cancels all outstanding polls for a given worker instance key.
func (h *Handler) CancelOutstandingWorkerPolls(ctx context.Context,
request *matchingservice.CancelOutstandingWorkerPollsRequest) (_ *matchingservice.CancelOutstandingWorkerPollsResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.CancelOutstandingWorkerPolls(ctx, request)
}
// CancelOutstandingWorkerPollsPartition cancels outstanding polls for a worker on a specific partition.
func (h *Handler) CancelOutstandingWorkerPollsPartition(ctx context.Context,
request *matchingservice.CancelOutstandingWorkerPollsPartitionRequest) (_ *matchingservice.CancelOutstandingWorkerPollsPartitionResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.CancelOutstandingWorkerPollsPartition(ctx, request)
}
// DescribeTaskQueue returns information about the target task queue, right now this API returns the
// pollers which polled this task queue in last few minutes. If includeTaskQueueStatus field is true,
// it will also return status of task queue's ackManager (readLevel, ackLevel, backlogCountHint and taskIDBlock).
func (h *Handler) DescribeTaskQueue(
ctx context.Context,
request *matchingservice.DescribeTaskQueueRequest,
defer log.CapturePanic(h.logger, &retError)
resp, err := h.engine.DescribeTaskQueue(ctx, request)
if err != nil {
return nil, err
}
// TODO: remove after 1.24.0-m3
if len(resp.DescResponse.Pollers) > 0 || resp.DescResponse.TaskQueueStatus != nil {
workflow_handler.go ×11
// Expand pollerinfo and task queue status into tags 1 and 2 for old frontend to handle
// proto incompatibility. This only works without ugly protowire code because
// workflowservice.DescribeTaskQueueResponse and the previous version of
// matchingservice.DescribeTaskQueueResponse have the same first two fields.
oldResp := &workflowservice.DescribeTaskQueueResponse{
Pollers: resp.DescResponse.Pollers,
TaskQueueStatus: resp.DescResponse.TaskQueueStatus,
}
if b, err := proto.Marshal(oldResp); err == nil {
resp.ProtoReflect().SetUnknown(protoreflect.RawFields(b))
}
}
}
func (h *Handler) DescribeVersionedTaskQueues(
ctx context.Context,
request *matchingservice.DescribeVersionedTaskQueuesRequest,
) (_ *matchingservice.DescribeVersionedTaskQueuesResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.DescribeVersionedTaskQueues(ctx, request)
}
// DescribeTaskQueuePartition returns information about the target task queue partition.
func (h *Handler) DescribeTaskQueuePartition(
ctx context.Context,
request *matchingservice.DescribeTaskQueuePartitionRequest,
) (_ *matchingservice.DescribeTaskQueuePartitionResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.DescribeTaskQueuePartition(ctx, request)
}
// ListTaskQueuePartitions returns information about partitions for a taskQueue
func (h *Handler) ListTaskQueuePartitions(
ctx context.Context,
request *matchingservice.ListTaskQueuePartitionsRequest,
) (_ *matchingservice.ListTaskQueuePartitionsResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.ListTaskQueuePartitions(ctx, request)
}
// UpdateWorkerVersioningRules allows updating the Build ID assignment and redirect rules for a given Task Queue.
func (h *Handler) UpdateWorkerVersioningRules(
ctx context.Context,
request *matchingservice.UpdateWorkerVersioningRulesRequest,
) (_ *matchingservice.UpdateWorkerVersioningRulesResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateWorkerVersioningRules(ctx, request)
}
// GetWorkerVersioningRules fetches the Build ID assignment and redirect rules for a Task Queue
func (h *Handler) GetWorkerVersioningRules(
ctx context.Context,
request *matchingservice.GetWorkerVersioningRulesRequest,
) (_ *matchingservice.GetWorkerVersioningRulesResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.GetWorkerVersioningRules(ctx, request)
}
// UpdateWorkerBuildIdCompatibility allows changing the worker versioning graph for a task queue
func (h *Handler) UpdateWorkerBuildIdCompatibility(
ctx context.Context,
request *matchingservice.UpdateWorkerBuildIdCompatibilityRequest,
) (_ *matchingservice.UpdateWorkerBuildIdCompatibilityResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateWorkerBuildIdCompatibility(ctx, request)
}
// GetWorkerBuildIdCompatibility fetches the worker versioning data for a task queue
func (h *Handler) GetWorkerBuildIdCompatibility(
ctx context.Context,
request *matchingservice.GetWorkerBuildIdCompatibilityRequest,
) (_ *matchingservice.GetWorkerBuildIdCompatibilityResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.GetWorkerBuildIdCompatibility(ctx, request)
}
func (h *Handler) GetTaskQueueUserData(
ctx context.Context,
request *matchingservice.GetTaskQueueUserDataRequest,
defer log.CapturePanic(h.logger, &retError)
return h.engine.GetTaskQueueUserData(ctx, request)
}
func (h *Handler) SyncDeploymentUserData(
ctx context.Context,
request *matchingservice.SyncDeploymentUserDataRequest,
) (_ *matchingservice.SyncDeploymentUserDataResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.SyncDeploymentUserData(ctx, request)
}
func (h *Handler) ApplyTaskQueueUserDataReplicationEvent(
ctx context.Context,
request *matchingservice.ApplyTaskQueueUserDataReplicationEventRequest,
) (_ *matchingservice.ApplyTaskQueueUserDataReplicationEventResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.ApplyTaskQueueUserDataReplicationEvent(ctx, request)
}
func (h *Handler) GetBuildIdTaskQueueMapping(
ctx context.Context,
request *matchingservice.GetBuildIdTaskQueueMappingRequest,
) (_ *matchingservice.GetBuildIdTaskQueueMappingResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.GetBuildIdTaskQueueMapping(ctx, request)
}
func (h *Handler) ForceUnloadTaskQueue(
ctx context.Context,
request *matchingservice.ForceUnloadTaskQueueRequest,
) (_ *matchingservice.ForceUnloadTaskQueueResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.ForceUnloadTaskQueue(ctx, request)
}
func (h *Handler) ForceUnloadTaskQueuePartition(
ctx context.Context,
request *matchingservice.ForceUnloadTaskQueuePartitionRequest,
) (_ *matchingservice.ForceUnloadTaskQueuePartitionResponse, retError error) {
service_grpc.pb.go ×19
defer log.CapturePanic(h.logger, &retError)
return h.engine.ForceUnloadTaskQueuePartition(ctx, request)
}
func (h *Handler) ForceLoadTaskQueuePartition(
ctx context.Context,
request *matchingservice.ForceLoadTaskQueuePartitionRequest,
) (_ *matchingservice.ForceLoadTaskQueuePartitionResponse, retError error) {
request_response.pb.go ×6
defer log.CapturePanic(h.logger, &retError)
return h.engine.ForceLoadTaskQueuePartition(ctx, request)
}
func (h *Handler) UpdateTaskQueueUserData(
ctx context.Context,
request *matchingservice.UpdateTaskQueueUserDataRequest,
) (_ *matchingservice.UpdateTaskQueueUserDataResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateTaskQueueUserData(ctx, request)
}
func (h *Handler) ReplicateTaskQueueUserData(
ctx context.Context,
request *matchingservice.ReplicateTaskQueueUserDataRequest,
) (_ *matchingservice.ReplicateTaskQueueUserDataResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.ReplicateTaskQueueUserData(ctx, request)
}
func (h *Handler) CheckTaskQueueUserDataPropagation(
ctx context.Context,
request *matchingservice.CheckTaskQueueUserDataPropagationRequest,
) (_ *matchingservice.CheckTaskQueueUserDataPropagationResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.CheckTaskQueueUserDataPropagation(ctx, request)
}
func (h *Handler) CheckTaskQueueVersionMembership(
ctx context.Context,
request *matchingservice.CheckTaskQueueVersionMembershipRequest,
) (_ *matchingservice.CheckTaskQueueVersionMembershipResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.CheckTaskQueueVersionMembership(ctx, request)
}
func (h *Handler) DispatchNexusTask(ctx context.Context, request *matchingservice.DispatchNexusTaskRequest) (_ *matchingservice.DispatchNexusTaskResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.DispatchNexusTask(ctx, request)
}
func (h *Handler) PollNexusTaskQueue(ctx context.Context, request *matchingservice.PollNexusTaskQueueRequest) (_ *matchingservice.PollNexusTaskQueueResponse, retError error) {
handler.go ×4
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetRequest().GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_NEXUS,
metrics.MatchingPollWorkflowTaskQueueScope,
)
// Only record on the initial handler call (ForwardedSource == ""), not on
// the forwarded call to the root partition, to avoid double-counting.
if request.GetForwardedSource() == "" {
h.recordNexusTaskRequest(ctx, request.GetNamespaceId(), request.GetRequest().GetTaskQueue().GetKind(), "PollNexusTaskQueue")
}
h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
}
ctx,
"PollNexusTaskQueue",
h.throttledLogger,
); err != nil {
return nil, err
}
}
func (h *Handler) RespondNexusTaskCompleted(ctx context.Context, request *matchingservice.RespondNexusTaskCompletedRequest) (_ *matchingservice.RespondNexusTaskCompletedResponse, retError error) {
request_response.pb.go ×1
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_NEXUS,
metrics.MatchingRespondNexusTaskCompletedScope,
)
h.recordNexusTaskRequest(ctx, request.GetNamespaceId(), request.GetTaskQueue().GetKind(), "RespondNexusTaskCompleted")
return h.engine.RespondNexusTaskCompleted(ctx, request, opMetrics)
}
func (h *Handler) RespondNexusTaskFailed(ctx context.Context, request *matchingservice.RespondNexusTaskFailedRequest) (_ *matchingservice.RespondNexusTaskFailedResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
opMetrics := h.opMetricsHandler(
request.GetNamespaceId(),
request.GetTaskQueue(),
enumspb.TASK_QUEUE_TYPE_NEXUS,
metrics.MatchingRespondNexusTaskFailedScope,
)
h.recordNexusTaskRequest(ctx, request.GetNamespaceId(), request.GetTaskQueue().GetKind(), "RespondNexusTaskFailed")
return h.engine.RespondNexusTaskFailed(ctx, request, opMetrics)
}
func (h *Handler) CreateNexusEndpoint(ctx context.Context, request *matchingservice.CreateNexusEndpointRequest) (_ *matchingservice.CreateNexusEndpointResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.CreateNexusEndpoint(ctx, request)
}
func (h *Handler) UpdateNexusEndpoint(ctx context.Context, request *matchingservice.UpdateNexusEndpointRequest) (_ *matchingservice.UpdateNexusEndpointResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateNexusEndpoint(ctx, request)
}
func (h *Handler) DeleteNexusEndpoint(ctx context.Context, request *matchingservice.DeleteNexusEndpointRequest) (_ *matchingservice.DeleteNexusEndpointResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.DeleteNexusEndpoint(ctx, request)
}
func (h *Handler) ListNexusEndpoints(ctx context.Context, request *matchingservice.ListNexusEndpointsRequest) (_ *matchingservice.ListNexusEndpointsResponse, retError error) {
service_grpc.pb.go ×20
defer log.CapturePanic(h.logger, &retError)
return h.engine.ListNexusEndpoints(ctx, request)
}
// RecordWorkerHeartbeat receive heartbeat request from the worker.
func (h *Handler) RecordWorkerHeartbeat(
ctx context.Context, request *matchingservice.RecordWorkerHeartbeatRequest,
defer log.CapturePanic(h.logger, &retError)
nsID := namespace.ID(request.GetNamespaceId())
nsName := h.namespaceName(nsID)
principal := headers.GetPrincipal(ctx)
h.workersRegistry.RecordWorkerHeartbeats(nsID, nsName, principal, request.GetHeartbeartRequest().GetWorkerHeartbeat())
return &matchingservice.RecordWorkerHeartbeatResponse{}, nil
}
// ListWorkers retrieves a list of workers in the specified namespace that match the provided filters.
func (h *Handler) ListWorkers(
_ context.Context, request *matchingservice.ListWorkersRequest,
) (_ *matchingservice.ListWorkersResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
nsID := namespace.ID(request.GetNamespaceId())
listRequest := request.GetListRequest()
resp, err := h.workersRegistry.ListWorkers(nsID, workers.ListWorkersParams{
Query: listRequest.GetQuery(),
PageSize: int(listRequest.GetPageSize()),
NextPageToken: listRequest.GetNextPageToken(),
IncludeSystemWorkers: listRequest.GetIncludeSystemWorkers(),
})
if err != nil {
return nil, err
}
// TODO: Stop populating workersInfo once all callers migrate to the Workers field.
workersInfo := make([]*workerpb.WorkerInfo, len(resp.Workers))
workersList := make([]*workerpb.WorkerListInfo, len(resp.Workers))
for i, heartbeat := range resp.Workers {
workersInfo[i] = &workerpb.WorkerInfo{
WorkerHeartbeat: heartbeat,
}
workersList[i] = workerHeartbeatToListInfo(heartbeat)
}
return &matchingservice.ListWorkersResponse{
WorkersInfo: workersInfo,
Workers: workersList,
NextPageToken: resp.NextPageToken,
}, nil
}
func workerHeartbeatToListInfo(hb *workerpb.WorkerHeartbeat) *workerpb.WorkerListInfo {
handler.go ×1
hostInfo := hb.GetHostInfo()
return &workerpb.WorkerListInfo{
WorkerInstanceKey: hb.GetWorkerInstanceKey(),
WorkerIdentity: hb.GetWorkerIdentity(),
TaskQueue: hb.GetTaskQueue(),
DeploymentVersion: hb.GetDeploymentVersion(),
SdkName: hb.GetSdkName(),
SdkVersion: hb.GetSdkVersion(),
Status: hb.GetStatus(),
StartTime: hb.GetStartTime(),
HostName: hostInfo.GetHostName(),
WorkerGroupingKey: hostInfo.GetWorkerGroupingKey(),
ProcessId: hostInfo.GetProcessId(),
Plugins: hb.GetPlugins(),
Drivers: hb.GetDrivers(),
}
}
func (h *Handler) CountWorkers(
_ context.Context, request *matchingservice.CountWorkersRequest,
) (_ *matchingservice.CountWorkersResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
nsID := namespace.ID(request.GetNamespaceId())
countRequest := request.GetCountRequest()
count, err := h.workersRegistry.CountWorkers(nsID, countRequest.GetQuery(), countRequest.GetIncludeSystemWorkers())
if err != nil {
return nil, err
}
return &matchingservice.CountWorkersResponse{
Count: count,
}, nil
}
func (h *Handler) UpdateFairnessState(
ctx context.Context, request *matchingservice.UpdateFairnessStateRequest,
) (_ *matchingservice.UpdateFairnessStateResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateFairnessState(ctx, request)
}
entry, err := h.namespaceRegistry.GetNamespaceByID(id)
if err != nil {
}
}
func (h *Handler) reportForwardedPerTaskQueueCounter(opMetrics metrics.Handler, namespaceId namespace.ID) {
request_response.pb.go ×6
metrics.ForwardedPerTaskQueueCounter.With(opMetrics).Record(1)
metrics.MatchingClientForwardedCounter.With(h.metricsHandler).
Record(
1,
metrics.OperationTag(metrics.MatchingAddWorkflowTaskScope),
metrics.NamespaceTag(h.namespaceName(namespaceId).String()),
metrics.ServiceRoleTag(metrics.MatchingRoleTagValue))
}
func (h *Handler) UpdateTaskQueueConfig(
ctx context.Context, request *matchingservice.UpdateTaskQueueConfigRequest,
) (_ *matchingservice.UpdateTaskQueueConfigResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
return h.engine.UpdateTaskQueueConfig(ctx, request)
}
func (h *Handler) DescribeWorker(
_ context.Context, request *matchingservice.DescribeWorkerRequest,
) (_ *matchingservice.DescribeWorkerResponse, retError error) {
defer log.CapturePanic(h.logger, &retError)
nsID := namespace.ID(request.GetNamespaceId())
hb, err := h.workersRegistry.DescribeWorker(
nsID, request.Request.GetWorkerInstanceKey())
if err != nil {
return nil, err
}
return &matchingservice.DescribeWorkerResponse{
WorkerInfo: &workerpb.WorkerInfo{
WorkerHeartbeat: hb,
},
}, nil
}