Admin UI add maintenance menu (#6944)
* add ui for maintenance * valid config loading. fix workers page. * refactor * grpc between admin and workers * add a long-running bidirectional grpc call between admin and worker * use the grpc call to heartbeat * use the grpc call to communicate * worker can remove the http client * admin uses http port + 10000 as its default grpc port * one task one package * handles connection failures gracefully with exponential backoff * grpc with insecure tls * grpc with optional tls * fix detecting tls * change time config from nano seconds to seconds * add tasks with 3 interfaces * compiles reducing hard coded * remove a couple of tasks * remove hard coded references * reduce hard coded values * remove hard coded values * remove hard coded from templ * refactor maintenance package * fix import cycle * simplify * simplify * auto register * auto register factory * auto register task types * self register types * refactor * simplify * remove one task * register ui * lazy init executor factories * use registered task types * DefaultWorkerConfig remove hard coded task types * remove more hard coded * implement get maintenance task * dynamic task configuration * "System Settings" should only have system level settings * adjust menu for tasks * ensure menu not collapsed * render job configuration well * use templ for ui of task configuration * fix ordering * fix bugs * saving duration in seconds * use value and unit for duration * Delete WORKER_REFACTORING_PLAN.md * Delete maintenance.json * Delete custom_worker_example.go * remove address from workers * remove old code from ec task * remove creating collection button * reconnect with exponential backoff * worker use security.toml * start admin server with tls info from security.toml * fix "weed admin" cli description
This commit is contained in:
@@ -13,6 +13,7 @@ gen:
|
||||
protoc mq_broker.proto --go_out=./mq_pb --go-grpc_out=./mq_pb --go_opt=paths=source_relative --go-grpc_opt=paths=source_relative
|
||||
protoc mq_schema.proto --go_out=./schema_pb --go-grpc_out=./schema_pb --go_opt=paths=source_relative --go-grpc_opt=paths=source_relative
|
||||
protoc mq_agent.proto --go_out=./mq_agent_pb --go-grpc_out=./mq_agent_pb --go_opt=paths=source_relative --go-grpc_opt=paths=source_relative
|
||||
protoc worker.proto --go_out=./worker_pb --go-grpc_out=./worker_pb --go_opt=paths=source_relative --go-grpc_opt=paths=source_relative
|
||||
# protoc filer.proto --java_out=../../other/java/client/src/main/java
|
||||
cp filer.proto ../../other/java/client/src/main/proto
|
||||
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/worker_pb"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -312,3 +313,10 @@ func WithOneOfGrpcFilerClients(streamingMode bool, filerAddresses []ServerAddres
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
func WithWorkerClient(streamingMode bool, workerAddress string, grpcDialOption grpc.DialOption, fn func(client worker_pb.WorkerServiceClient) error) error {
|
||||
return WithGrpcClient(streamingMode, 0, func(grpcConnection *grpc.ClientConn) error {
|
||||
client := worker_pb.NewWorkerServiceClient(grpcConnection)
|
||||
return fn(client)
|
||||
}, workerAddress, false, grpcDialOption)
|
||||
}
|
||||
|
||||
142
weed/pb/worker.proto
Normal file
142
weed/pb/worker.proto
Normal file
@@ -0,0 +1,142 @@
|
||||
syntax = "proto3";
|
||||
|
||||
package worker_pb;
|
||||
|
||||
option go_package = "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb";
|
||||
|
||||
// WorkerService provides bidirectional communication between admin and worker
|
||||
service WorkerService {
|
||||
// WorkerStream maintains a bidirectional stream for worker communication
|
||||
rpc WorkerStream(stream WorkerMessage) returns (stream AdminMessage);
|
||||
}
|
||||
|
||||
// WorkerMessage represents messages from worker to admin
|
||||
message WorkerMessage {
|
||||
string worker_id = 1;
|
||||
int64 timestamp = 2;
|
||||
|
||||
oneof message {
|
||||
WorkerRegistration registration = 3;
|
||||
WorkerHeartbeat heartbeat = 4;
|
||||
TaskRequest task_request = 5;
|
||||
TaskUpdate task_update = 6;
|
||||
TaskComplete task_complete = 7;
|
||||
WorkerShutdown shutdown = 8;
|
||||
}
|
||||
}
|
||||
|
||||
// AdminMessage represents messages from admin to worker
|
||||
message AdminMessage {
|
||||
string admin_id = 1;
|
||||
int64 timestamp = 2;
|
||||
|
||||
oneof message {
|
||||
RegistrationResponse registration_response = 3;
|
||||
HeartbeatResponse heartbeat_response = 4;
|
||||
TaskAssignment task_assignment = 5;
|
||||
TaskCancellation task_cancellation = 6;
|
||||
AdminShutdown admin_shutdown = 7;
|
||||
}
|
||||
}
|
||||
|
||||
// WorkerRegistration message when worker connects
|
||||
message WorkerRegistration {
|
||||
string worker_id = 1;
|
||||
string address = 2;
|
||||
repeated string capabilities = 3;
|
||||
int32 max_concurrent = 4;
|
||||
map<string, string> metadata = 5;
|
||||
}
|
||||
|
||||
// RegistrationResponse confirms worker registration
|
||||
message RegistrationResponse {
|
||||
bool success = 1;
|
||||
string message = 2;
|
||||
string assigned_worker_id = 3;
|
||||
}
|
||||
|
||||
// WorkerHeartbeat sent periodically by worker
|
||||
message WorkerHeartbeat {
|
||||
string worker_id = 1;
|
||||
string status = 2;
|
||||
int32 current_load = 3;
|
||||
int32 max_concurrent = 4;
|
||||
repeated string current_task_ids = 5;
|
||||
int32 tasks_completed = 6;
|
||||
int32 tasks_failed = 7;
|
||||
int64 uptime_seconds = 8;
|
||||
}
|
||||
|
||||
// HeartbeatResponse acknowledges heartbeat
|
||||
message HeartbeatResponse {
|
||||
bool success = 1;
|
||||
string message = 2;
|
||||
}
|
||||
|
||||
// TaskRequest from worker asking for new tasks
|
||||
message TaskRequest {
|
||||
string worker_id = 1;
|
||||
repeated string capabilities = 2;
|
||||
int32 available_slots = 3;
|
||||
}
|
||||
|
||||
// TaskAssignment from admin to worker
|
||||
message TaskAssignment {
|
||||
string task_id = 1;
|
||||
string task_type = 2;
|
||||
TaskParams params = 3;
|
||||
int32 priority = 4;
|
||||
int64 created_time = 5;
|
||||
map<string, string> metadata = 6;
|
||||
}
|
||||
|
||||
// TaskParams contains task-specific parameters
|
||||
message TaskParams {
|
||||
uint32 volume_id = 1;
|
||||
string server = 2;
|
||||
string collection = 3;
|
||||
string data_center = 4;
|
||||
string rack = 5;
|
||||
repeated string replicas = 6;
|
||||
map<string, string> parameters = 7;
|
||||
}
|
||||
|
||||
// TaskUpdate reports task progress
|
||||
message TaskUpdate {
|
||||
string task_id = 1;
|
||||
string worker_id = 2;
|
||||
string status = 3;
|
||||
float progress = 4;
|
||||
string message = 5;
|
||||
map<string, string> metadata = 6;
|
||||
}
|
||||
|
||||
// TaskComplete reports task completion
|
||||
message TaskComplete {
|
||||
string task_id = 1;
|
||||
string worker_id = 2;
|
||||
bool success = 3;
|
||||
string error_message = 4;
|
||||
int64 completion_time = 5;
|
||||
map<string, string> result_metadata = 6;
|
||||
}
|
||||
|
||||
// TaskCancellation from admin to cancel a task
|
||||
message TaskCancellation {
|
||||
string task_id = 1;
|
||||
string reason = 2;
|
||||
bool force = 3;
|
||||
}
|
||||
|
||||
// WorkerShutdown notifies admin that worker is shutting down
|
||||
message WorkerShutdown {
|
||||
string worker_id = 1;
|
||||
string reason = 2;
|
||||
repeated string pending_task_ids = 3;
|
||||
}
|
||||
|
||||
// AdminShutdown notifies worker that admin is shutting down
|
||||
message AdminShutdown {
|
||||
string reason = 1;
|
||||
int32 graceful_shutdown_seconds = 2;
|
||||
}
|
||||
1724
weed/pb/worker_pb/worker.pb.go
Normal file
1724
weed/pb/worker_pb/worker.pb.go
Normal file
File diff suppressed because it is too large
Load Diff
121
weed/pb/worker_pb/worker_grpc.pb.go
Normal file
121
weed/pb/worker_pb/worker_grpc.pb.go
Normal file
@@ -0,0 +1,121 @@
|
||||
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
|
||||
// versions:
|
||||
// - protoc-gen-go-grpc v1.5.1
|
||||
// - protoc v5.29.3
|
||||
// source: worker.proto
|
||||
|
||||
package worker_pb
|
||||
|
||||
import (
|
||||
context "context"
|
||||
grpc "google.golang.org/grpc"
|
||||
codes "google.golang.org/grpc/codes"
|
||||
status "google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// This is a compile-time assertion to ensure that this generated file
|
||||
// is compatible with the grpc package it is being compiled against.
|
||||
// Requires gRPC-Go v1.64.0 or later.
|
||||
const _ = grpc.SupportPackageIsVersion9
|
||||
|
||||
const (
|
||||
WorkerService_WorkerStream_FullMethodName = "/worker_pb.WorkerService/WorkerStream"
|
||||
)
|
||||
|
||||
// WorkerServiceClient is the client API for WorkerService 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.
|
||||
//
|
||||
// WorkerService provides bidirectional communication between admin and worker
|
||||
type WorkerServiceClient interface {
|
||||
// WorkerStream maintains a bidirectional stream for worker communication
|
||||
WorkerStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[WorkerMessage, AdminMessage], error)
|
||||
}
|
||||
|
||||
type workerServiceClient struct {
|
||||
cc grpc.ClientConnInterface
|
||||
}
|
||||
|
||||
func NewWorkerServiceClient(cc grpc.ClientConnInterface) WorkerServiceClient {
|
||||
return &workerServiceClient{cc}
|
||||
}
|
||||
|
||||
func (c *workerServiceClient) WorkerStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[WorkerMessage, AdminMessage], error) {
|
||||
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
|
||||
stream, err := c.cc.NewStream(ctx, &WorkerService_ServiceDesc.Streams[0], WorkerService_WorkerStream_FullMethodName, cOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
x := &grpc.GenericClientStream[WorkerMessage, AdminMessage]{ClientStream: stream}
|
||||
return x, nil
|
||||
}
|
||||
|
||||
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
|
||||
type WorkerService_WorkerStreamClient = grpc.BidiStreamingClient[WorkerMessage, AdminMessage]
|
||||
|
||||
// WorkerServiceServer is the server API for WorkerService service.
|
||||
// All implementations must embed UnimplementedWorkerServiceServer
|
||||
// for forward compatibility.
|
||||
//
|
||||
// WorkerService provides bidirectional communication between admin and worker
|
||||
type WorkerServiceServer interface {
|
||||
// WorkerStream maintains a bidirectional stream for worker communication
|
||||
WorkerStream(grpc.BidiStreamingServer[WorkerMessage, AdminMessage]) error
|
||||
mustEmbedUnimplementedWorkerServiceServer()
|
||||
}
|
||||
|
||||
// UnimplementedWorkerServiceServer must be embedded to have
|
||||
// forward compatible implementations.
|
||||
//
|
||||
// NOTE: this should be embedded by value instead of pointer to avoid a nil
|
||||
// pointer dereference when methods are called.
|
||||
type UnimplementedWorkerServiceServer struct{}
|
||||
|
||||
func (UnimplementedWorkerServiceServer) WorkerStream(grpc.BidiStreamingServer[WorkerMessage, AdminMessage]) error {
|
||||
return status.Errorf(codes.Unimplemented, "method WorkerStream not implemented")
|
||||
}
|
||||
func (UnimplementedWorkerServiceServer) mustEmbedUnimplementedWorkerServiceServer() {}
|
||||
func (UnimplementedWorkerServiceServer) testEmbeddedByValue() {}
|
||||
|
||||
// UnsafeWorkerServiceServer may be embedded to opt out of forward compatibility for this service.
|
||||
// Use of this interface is not recommended, as added methods to WorkerServiceServer will
|
||||
// result in compilation errors.
|
||||
type UnsafeWorkerServiceServer interface {
|
||||
mustEmbedUnimplementedWorkerServiceServer()
|
||||
}
|
||||
|
||||
func RegisterWorkerServiceServer(s grpc.ServiceRegistrar, srv WorkerServiceServer) {
|
||||
// If the following call pancis, it indicates UnimplementedWorkerServiceServer was
|
||||
// embedded by pointer and is nil. This will cause panics if an
|
||||
// unimplemented method is ever invoked, so we test this at initialization
|
||||
// time to prevent it from happening at runtime later due to I/O.
|
||||
if t, ok := srv.(interface{ testEmbeddedByValue() }); ok {
|
||||
t.testEmbeddedByValue()
|
||||
}
|
||||
s.RegisterService(&WorkerService_ServiceDesc, srv)
|
||||
}
|
||||
|
||||
func _WorkerService_WorkerStream_Handler(srv interface{}, stream grpc.ServerStream) error {
|
||||
return srv.(WorkerServiceServer).WorkerStream(&grpc.GenericServerStream[WorkerMessage, AdminMessage]{ServerStream: stream})
|
||||
}
|
||||
|
||||
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
|
||||
type WorkerService_WorkerStreamServer = grpc.BidiStreamingServer[WorkerMessage, AdminMessage]
|
||||
|
||||
// WorkerService_ServiceDesc is the grpc.ServiceDesc for WorkerService service.
|
||||
// It's only intended for direct use with grpc.RegisterService,
|
||||
// and not to be introspected or modified (even as a copy)
|
||||
var WorkerService_ServiceDesc = grpc.ServiceDesc{
|
||||
ServiceName: "worker_pb.WorkerService",
|
||||
HandlerType: (*WorkerServiceServer)(nil),
|
||||
Methods: []grpc.MethodDesc{},
|
||||
Streams: []grpc.StreamDesc{
|
||||
{
|
||||
StreamName: "WorkerStream",
|
||||
Handler: _WorkerService_WorkerStream_Handler,
|
||||
ServerStreams: true,
|
||||
ClientStreams: true,
|
||||
},
|
||||
},
|
||||
Metadata: "worker.proto",
|
||||
}
|
||||
Reference in New Issue
Block a user