tdbg for list fairness tasks (#8010)

## What changed?
Added new flag to tdbg's `list-tasks` command to query the fairness
table instead.

## Why?
Allow operation on fairness table, too.

## How did you test it?
- [ ] built
- [x] run locally and tested manually
- [ ] covered by existing tests
- [ ] added new unit test(s)
- [ ] added new functional test(s)

```
./tdbg taskqueue list-tasks --task-queue-type TASK_QUEUE_TYPE_WORKFLOW --task-queue test --min-pass 1
```

It's late ... I didn't manage to see a result, though. But it didn't
fail and I saw the query being executed in the debugger.
This commit is contained in:
Stephan Behnke
2025-07-07 18:06:56 -07:00
committed by David Reiss
parent e7d17e0b63
commit b157875c3e
13 changed files with 88 additions and 17 deletions

View File

@@ -3285,6 +3285,7 @@ type GetTaskQueueTasksRequest struct {
Namespace string `protobuf:"bytes,1,opt,name=namespace,proto3" json:"namespace,omitempty"`
TaskQueue string `protobuf:"bytes,2,opt,name=task_queue,json=taskQueue,proto3" json:"task_queue,omitempty"`
TaskQueueType v16.TaskQueueType `protobuf:"varint,3,opt,name=task_queue_type,json=taskQueueType,proto3,enum=temporal.api.enums.v1.TaskQueueType" json:"task_queue_type,omitempty"`
MinPass int64 `protobuf:"varint,9,opt,name=min_pass,json=minPass,proto3" json:"min_pass,omitempty"`
MinTaskId int64 `protobuf:"varint,4,opt,name=min_task_id,json=minTaskId,proto3" json:"min_task_id,omitempty"`
MaxTaskId int64 `protobuf:"varint,5,opt,name=max_task_id,json=maxTaskId,proto3" json:"max_task_id,omitempty"`
BatchSize int32 `protobuf:"varint,6,opt,name=batch_size,json=batchSize,proto3" json:"batch_size,omitempty"`
@@ -3345,6 +3346,13 @@ func (x *GetTaskQueueTasksRequest) GetTaskQueueType() v16.TaskQueueType {
return v16.TaskQueueType(0)
}
func (x *GetTaskQueueTasksRequest) GetMinPass() int64 {
if x != nil {
return x.MinPass
}
return 0
}
func (x *GetTaskQueueTasksRequest) GetMinTaskId() int64 {
if x != nil {
return x.MinTaskId
@@ -5593,12 +5601,13 @@ const file_temporal_server_api_adminservice_v1_request_response_proto_rawDesc =
"endEventId\x12\x1f\n" +
"\vend_version\x18\b \x01(\x03R\n" +
"endVersion\" \n" +
"\x1eResendReplicationTasksResponse\"\xc8\x02\n" +
"\x1eResendReplicationTasksResponse\"\xe3\x02\n" +
"\x18GetTaskQueueTasksRequest\x12\x1c\n" +
"\tnamespace\x18\x01 \x01(\tR\tnamespace\x12\x1d\n" +
"\n" +
"task_queue\x18\x02 \x01(\tR\ttaskQueue\x12L\n" +
"\x0ftask_queue_type\x18\x03 \x01(\x0e2$.temporal.api.enums.v1.TaskQueueTypeR\rtaskQueueType\x12\x1e\n" +
"\x0ftask_queue_type\x18\x03 \x01(\x0e2$.temporal.api.enums.v1.TaskQueueTypeR\rtaskQueueType\x12\x19\n" +
"\bmin_pass\x18\t \x01(\x03R\aminPass\x12\x1e\n" +
"\vmin_task_id\x18\x04 \x01(\x03R\tminTaskId\x12\x1e\n" +
"\vmax_task_id\x18\x05 \x01(\x03R\tmaxTaskId\x12\x1d\n" +
"\n" +

View File

@@ -346,7 +346,7 @@ func (d *matchingTaskStoreV1) GetTasks(
request *p.GetTasksRequest,
) (*p.InternalGetTasksResponse, error) {
if request.InclusiveMinPass != 0 {
return nil, serviceerror.NewInternal("invalid GetTasks request on queue")
return nil, serviceerror.NewInternal("invalid GetTasks request on queue: InclusiveMinPass is not supported")
}
// Reading taskqueue tasks need to be quorum level consistent, otherwise we could lose tasks

View File

@@ -287,9 +287,12 @@ func (d *matchingTaskStoreV2) GetTasks(
ctx context.Context,
request *p.GetTasksRequest,
) (*p.InternalGetTasksResponse, error) {
// Require starting from pass 1.
if request.InclusiveMinPass < 1 || request.ExclusiveMaxTaskID != math.MaxInt64 {
return nil, serviceerror.NewInternal("invalid GetTasks request on fair queue")
if request.InclusiveMinPass < 1 {
return nil, serviceerror.NewInternal("invalid GetTasks request on fair queue: InclusiveMinPass must be >= 1")
}
if request.ExclusiveMaxTaskID != math.MaxInt64 {
// ExclusiveMaxTaskID is not supported in fair queue.
return nil, serviceerror.NewInternal("invalid GetTasks request on fair queue: ExclusiveMaxTaskID is not supported")
}
// Reading taskqueue tasks need to be quorum level consistent, otherwise we could lose tasks

View File

@@ -377,6 +377,7 @@ message GetTaskQueueTasksRequest {
string namespace = 1;
string task_queue = 2;
temporal.api.enums.v1.TaskQueueType task_queue_type = 3;
int64 min_pass = 9;
int64 min_task_id = 4;
int64 max_task_id = 5;
int32 batch_size = 6;

View File

@@ -0,0 +1,6 @@
{
"CurrVersion": "1.13",
"MinCompatibleVersion": "1.0",
"Description": "Adds tasks_v2 table for fairness tasks",
"SchemaUpdateCqlFiles": ["tasks_v2.cql"]
}

View File

@@ -0,0 +1,18 @@
CREATE TABLE tasks_v2
(
namespace_id uuid,
task_queue_name text,
task_queue_type int, -- enum TaskQueueType {ActivityTask, WorkflowTask}
type int, -- enum rowType {Task, TaskQueue} and subqueue id
pass bigint, -- pass for tasks (see stride scheduling algorithm for fairness)
task_id bigint, -- unique identifier for tasks
range_id bigint, -- used to ensure that only one process can write to the table
ack_level bigint, -- ack level for the task queue
task blob,
task_encoding text,
task_queue blob,
task_queue_encoding text,
PRIMARY KEY ((namespace_id, task_queue_name, task_queue_type), type, pass, task_id)
) WITH COMPACTION = {
'class': 'org.apache.cassandra.db.compaction.LeveledCompactionStrategy'
};

View File

@@ -3,4 +3,4 @@ package cassandra
// NOTE: whenever there is a new database schema update, plz update the following versions
// Version is the Cassandra database release version
const Version = "1.12"
const Version = "1.13"

View File

@@ -6,6 +6,7 @@ import (
"fmt"
"io"
"maps"
"math"
"net"
"strings"
"sync/atomic"
@@ -85,6 +86,7 @@ type (
persistenceExecutionName string
namespaceReplicationQueue persistence.NamespaceReplicationQueue
taskManager persistence.TaskManager
fairTaskManager persistence.FairTaskManager
clusterMetadataManager persistence.ClusterMetadataManager
persistenceMetadataManager persistence.MetadataManager
clientFactory serverClient.Factory
@@ -117,6 +119,7 @@ type (
visibilityMgr manager.VisibilityManager
Logger log.Logger
TaskManager persistence.TaskManager
FairTaskManager persistence.FairTaskManager
PersistenceExecutionManager persistence.ExecutionManager
ClusterMetadataManager persistence.ClusterMetadataManager
PersistenceMetadataManager persistence.MetadataManager
@@ -189,6 +192,7 @@ func NewAdminHandler(
persistenceExecutionName: args.PersistenceExecutionManager.GetName(),
namespaceReplicationQueue: args.NamespaceReplicationQueue,
taskManager: args.TaskManager,
fairTaskManager: args.FairTaskManager,
clusterMetadataManager: args.ClusterMetadataManager,
persistenceMetadataManager: args.PersistenceMetadataManager,
clientFactory: args.ClientFactory,
@@ -1563,12 +1567,24 @@ func (adh *AdminHandler) GetTaskQueueTasks(
return nil, err
}
resp, err := adh.taskManager.GetTasks(ctx, &persistence.GetTasksRequest{
var taskManager persistence.TaskManager
if request.GetMinPass() != 0 {
if adh.fairTaskManager == nil {
return nil, serviceerror.NewInvalidArgument("Fairness table is not available on this cluster")
}
taskManager = adh.fairTaskManager
request.MaxTaskId = math.MaxInt64 // required for fairness GetTasks call
} else {
taskManager = adh.taskManager
}
resp, err := taskManager.GetTasks(ctx, &persistence.GetTasksRequest{
NamespaceID: namespaceID.String(),
TaskQueue: request.GetTaskQueue(),
TaskType: request.GetTaskQueueType(),
InclusiveMinTaskID: request.GetMinTaskId(),
ExclusiveMaxTaskID: request.GetMaxTaskId(),
InclusiveMinPass: request.GetMinPass(),
Subqueue: int(request.GetSubqueue()),
PageSize: int(request.GetBatchSize()),
NextPageToken: request.NextPageToken,

View File

@@ -153,6 +153,7 @@ func (s *adminHandlerSuite) SetupTest() {
s.mockVisibilityMgr,
s.mockResource.GetLogger(),
s.mockResource.GetTaskManager(),
s.mockResource.GetTaskManager(),
s.mockResource.GetExecutionManager(),
s.mockResource.GetClusterMetadataManager(),
s.mockResource.GetMetadataManager(),

View File

@@ -601,6 +601,7 @@ func AdminHandlerProvider(
logger log.SnTaggedLogger,
namespaceReplicationQueue persistence.NamespaceReplicationQueue,
taskManager persistence.TaskManager,
fairTaskManager persistence.FairTaskManager,
persistenceExecutionManager persistence.ExecutionManager,
clusterMetadataManager persistence.ClusterMetadataManager,
persistenceMetadataManager persistence.MetadataManager,
@@ -630,6 +631,7 @@ func AdminHandlerProvider(
visibilityMgr,
logger,
taskManager,
fairTaskManager,
persistenceExecutionManager,
clusterMetadataManager,
persistenceMetadataManager,

View File

@@ -68,4 +68,6 @@ var (
FlagBuildIDs = "select-build-id"
FlagUnversioned = "select-unversioned"
FlagAllActive = "select-all-active"
FlagFair = "fair"
FlagMinPass = "min-pass"
)

View File

@@ -28,14 +28,16 @@ func AdminListTaskQueueTasks(c *cli.Context, clientFactory ClientFactory) error
}
minTaskID := c.Int64(FlagMinTaskID)
maxTaskID := c.Int64(FlagMaxTaskID)
pageSize := defaultPageSize
if c.IsSet(FlagPageSize) {
pageSize = c.Int(FlagPageSize)
}
pageSize := c.Int(FlagPageSize)
workflowID := c.String(FlagWorkflowID)
runID := c.String(FlagRunID)
subqueue := c.Int(FlagSubqueue)
var minPass int64
if c.Bool(FlagFair) {
minPass = c.Int64(FlagMinPass)
} else if c.IsSet(FlagMinPass) {
return fmt.Errorf("flag --%s is only valid with --%s", FlagMinPass, FlagFair)
}
client := clientFactory.AdminClient(c)
req := &adminservice.GetTaskQueueTasksRequest{
@@ -46,11 +48,13 @@ func AdminListTaskQueueTasks(c *cli.Context, clientFactory ClientFactory) error
MaxTaskId: maxTaskID,
BatchSize: int32(pageSize),
Subqueue: int32(subqueue),
MinPass: minPass,
}
ctx, cancel := newContext(c)
defer cancel()
paginationFunc := func(paginationToken []byte) ([]interface{}, []byte, error) {
ctx, cancel := newContext(c)
defer cancel()
req.NextPageToken = paginationToken
response, err := client.GetTaskQueueTasks(ctx, req)
if err != nil {
@@ -78,7 +82,7 @@ func AdminListTaskQueueTasks(c *cli.Context, clientFactory ClientFactory) error
for _, task := range tasks {
items = append(items, task)
}
return items, nil, nil
return items, response.NextPageToken, nil
}
if err := paginate(c, paginationFunc, pageSize); err != nil {

View File

@@ -462,7 +462,7 @@ func newAdminTaskQueueCommands(clientFactory ClientFactory) []*cli.Command {
return []*cli.Command{
{
Name: "list-tasks",
Usage: "List tasks of a task queue",
Usage: "List tasks of a task queue. Use --fair to list fairness tasks.",
Flags: []cli.Flag{
&cli.BoolFlag{
Name: FlagMore,
@@ -500,6 +500,15 @@ func newAdminTaskQueueCommands(clientFactory ClientFactory) []*cli.Command {
Name: FlagPrintJSON,
Usage: "Print in raw json format",
},
&cli.BoolFlag{
Name: FlagFair,
Usage: "Query fairness tasks",
},
&cli.Int64Flag{
Name: FlagMinPass,
Usage: "Minimum pass (fairness task only)",
Value: 1,
},
},
Action: func(c *cli.Context) error {
return AdminListTaskQueueTasks(c, clientFactory)