feat: Implement plugin system infrastructure

This commit is contained in:
Chris Lu committed 2026-02-17 00:38:07 -08:00
1 parent 3bd20e6a10
commit 7f9e93384a
11 files changed
+8384

No files matched your search

+307
View File
@@ -0,0 +1,307 @@
package plugin
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
)
// ConfigManager handles configuration persistence and retrieval
type ConfigManager struct {
dataDir string
configs map[string]*JobTypeConfig
mu sync.RWMutex
}
// NewConfigManager creates a new configuration manager
func NewConfigManager(dataDir string) (*ConfigManager, error) {
cm := &ConfigManager{
dataDir: filepath.Join(dataDir, "plugins"),
configs: make(map[string]*JobTypeConfig),
}
// Create plugins directory if it doesn't exist
if err := os.MkdirAll(cm.dataDir, 0755); err != nil {
return nil, fmt.Errorf("failed to create plugin config directory: %w", err)
}
return cm, nil
}
// LoadConfig loads a configuration from disk
func (cm *ConfigManager) LoadConfig(jobType string) (*JobTypeConfig, error) {
cm.mu.Lock()
defer cm.mu.Unlock()
// Check cache first
if config, exists := cm.configs[jobType]; exists {
return config, nil
}
configPath := filepath.Join(cm.dataDir, jobType+".json")
// If file doesn't exist, create default config
if _, err := os.Stat(configPath); err != nil {
if os.IsNotExist(err) {
config := &JobTypeConfig{
JobType: jobType,
Enabled: false,
AdminConfig: make([]*plugin_pb.ConfigFieldValue, 0),
WorkerConfig: make([]*plugin_pb.ConfigFieldValue, 0),
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
CreatedBy: "system",
}
cm.configs[jobType] = config
return config, nil
}
return nil, fmt.Errorf("failed to access config file: %w", err)
}
// Load from file
data, err := os.ReadFile(configPath)
if err != nil {
return nil, fmt.Errorf("failed to read config file: %w", err)
}
var configData map[string]interface{}
if err := json.Unmarshal(data, &configData); err != nil {
return nil, fmt.Errorf("failed to parse config file: %w", err)
}
config := &JobTypeConfig{
JobType: jobType,
AdminConfig: make([]*plugin_pb.ConfigFieldValue, 0),
WorkerConfig: make([]*plugin_pb.ConfigFieldValue, 0),
}
// Parse fields
if enabled, ok := configData["enabled"].(bool); ok {
config.Enabled = enabled
}
if createdBy, ok := configData["created_by"].(string); ok {
config.CreatedBy = createdBy
}
if createdAt, ok := configData["created_at"].(string); ok {
if t, err := time.Parse(time.RFC3339, createdAt); err == nil {
config.CreatedAt = t
}
}
if updatedAt, ok := configData["updated_at"].(string); ok {
if t, err := time.Parse(time.RFC3339, updatedAt); err == nil {
config.UpdatedAt = t
}
}
cm.configs[jobType] = config
return config, nil
}
// SaveConfig saves a configuration to disk
func (cm *ConfigManager) SaveConfig(config *JobTypeConfig) error {
cm.mu.Lock()
defer cm.mu.Unlock()
config.UpdatedAt = time.Now()
cm.configs[config.JobType] = config
configPath := filepath.Join(cm.dataDir, config.JobType+".json")
// Convert to JSON-serializable format
configData := map[string]interface{}{
"job_type": config.JobType,
"enabled": config.Enabled,
"created_by": config.CreatedBy,
"created_at": config.CreatedAt.Format(time.RFC3339),
"updated_at": config.UpdatedAt.Format(time.RFC3339),
"admin_config": cm.configValuesToJSON(config.AdminConfig),
"worker_config": cm.configValuesToJSON(config.WorkerConfig),
}
data, err := json.MarshalIndent(configData, "", " ")
if err != nil {
return fmt.Errorf("failed to marshal config: %w", err)
}
if err := os.WriteFile(configPath, data, 0644); err != nil {
return fmt.Errorf("failed to write config file: %w", err)
}
return nil
}
// GetConfig retrieves a configuration
func (cm *ConfigManager) GetConfig(jobType string) (*JobTypeConfig, error) {
cm.mu.RLock()
if config, exists := cm.configs[jobType]; exists {
cm.mu.RUnlock()
return config, nil
}
cm.mu.RUnlock()
// Load from disk if not in cache
return cm.LoadConfig(jobType)
}
// UpdateConfig updates a configuration
func (cm *ConfigManager) UpdateConfig(jobType string, enabled bool, adminConfig, workerConfig []*plugin_pb.ConfigFieldValue, updatedBy string) (*JobTypeConfig, error) {
config, err := cm.GetConfig(jobType)
if err != nil {
return nil, err
}
config.SetConfig(enabled, adminConfig, workerConfig)
config.CreatedBy = updatedBy
if err := cm.SaveConfig(config); err != nil {
return nil, err
}
return config, nil
}
// ListConfigs lists all configurations
func (cm *ConfigManager) ListConfigs() []*JobTypeConfig {
cm.mu.RLock()
defer cm.mu.RUnlock()
configs := make([]*JobTypeConfig, 0, len(cm.configs))
for _, config := range cm.configs {
configs = append(configs, config)
}
return configs
}
// DeleteConfig deletes a configuration
func (cm *ConfigManager) DeleteConfig(jobType string) error {
cm.mu.Lock()
defer cm.mu.Unlock()
configPath := filepath.Join(cm.dataDir, jobType+".json")
if err := os.Remove(configPath); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("failed to delete config file: %w", err)
}
delete(cm.configs, jobType)
return nil
}
// configValuesToJSON converts ConfigFieldValue slice to JSON
func (cm *ConfigManager) configValuesToJSON(values []*plugin_pb.ConfigFieldValue) []map[string]interface{} {
result := make([]map[string]interface{}, 0, len(values))
for _, val := range values {
item := map[string]interface{}{
"field_name": val.FieldName,
}
// Include non-empty values
if val.StringValue != "" {
item["string_value"] = val.StringValue
}
if val.IntValue != 0 {
item["int_value"] = val.IntValue
}
if val.FloatValue != 0 {
item["float_value"] = val.FloatValue
}
if val.BoolValue {
item["bool_value"] = val.BoolValue
}
if val.DurationValue != nil {
item["duration_value"] = val.DurationValue.String()
}
if len(val.MultiselectValues) > 0 {
item["multiselect_values"] = val.MultiselectValues
}
if val.JsonValue != "" {
item["json_value"] = val.JsonValue
}
result = append(result, item)
}
return result
}
// LoadAllConfigs loads all configurations from disk
func (cm *ConfigManager) LoadAllConfigs() error {
cm.mu.Lock()
defer cm.mu.Unlock()
entries, err := os.ReadDir(cm.dataDir)
if err != nil {
return fmt.Errorf("failed to read config directory: %w", err)
}
for _, entry := range entries {
if entry.IsDir() {
continue
}
// Only load .json files
if filepath.Ext(entry.Name()) != ".json" {
continue
}
jobType := entry.Name()[:len(entry.Name())-5] // Remove .json
configPath := filepath.Join(cm.dataDir, entry.Name())
data, err := os.ReadFile(configPath)
if err != nil {
continue
}
var configData map[string]interface{}
if err := json.Unmarshal(data, &configData); err != nil {
continue
}
config := &JobTypeConfig{
JobType: jobType,
AdminConfig: make([]*plugin_pb.ConfigFieldValue, 0),
WorkerConfig: make([]*plugin_pb.ConfigFieldValue, 0),
}
// Parse fields
if enabled, ok := configData["enabled"].(bool); ok {
config.Enabled = enabled
}
if createdBy, ok := configData["created_by"].(string); ok {
config.CreatedBy = createdBy
}
if createdAt, ok := configData["created_at"].(string); ok {
if t, err := time.Parse(time.RFC3339, createdAt); err == nil {
config.CreatedAt = t
}
}
if updatedAt, ok := configData["updated_at"].(string); ok {
if t, err := time.Parse(time.RFC3339, updatedAt); err == nil {
config.UpdatedAt = t
}
}
cm.configs[jobType] = config
}
return nil
}
// GetDataDir returns the data directory path
func (cm *ConfigManager) GetDataDir() string {
return cm.dataDir
}
+373
View File
@@ -0,0 +1,373 @@
package plugin
import (
"context"
"fmt"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
"google.golang.org/grpc"
)
// Dispatcher handles job detection and dispatch to executors
type Dispatcher struct {
registry *Registry
queues map[string]*JobQueue
configMgr *ConfigManager
mu sync.RWMutex
// Detection scheduling
detectionIntervals map[string]time.Duration
detectionTickers map[string]*time.Ticker
stopChan chan struct{}
// Job counter for ID generation
jobCounter uint64
jobCounterMu sync.Mutex
// Active streams for execution
executionStreams map[string]*ExecutionStream
streamsMu sync.RWMutex
}
// ExecutionStream tracks an active job execution stream
type ExecutionStream struct {
JobID string
PluginID string
Connection *grpc.ClientConn
Context context.Context
Cancel context.CancelFunc
StartTime time.Time
LastUpdate time.Time
}
// NewDispatcher creates a new job dispatcher
func NewDispatcher(registry *Registry, configMgr *ConfigManager) *Dispatcher {
return &Dispatcher{
registry: registry,
queues: make(map[string]*JobQueue),
configMgr: configMgr,
detectionIntervals: make(map[string]time.Duration),
detectionTickers: make(map[string]*time.Ticker),
stopChan: make(chan struct{}),
executionStreams: make(map[string]*ExecutionStream),
}
}
// RegisterJobType registers a job type for detection and dispatch
func (d *Dispatcher) RegisterJobType(jobType string, detectionInterval time.Duration) error {
d.mu.Lock()
defer d.mu.Unlock()
if _, exists := d.queues[jobType]; exists {
return fmt.Errorf("job type %s already registered", jobType)
}
d.queues[jobType] = NewJobQueue(jobType)
d.detectionIntervals[jobType] = detectionInterval
return nil
}
// Start starts the dispatcher with detection scheduling
func (d *Dispatcher) Start() {
d.mu.Lock()
defer d.mu.Unlock()
for jobType, interval := range d.detectionIntervals {
ticker := time.NewTicker(interval)
d.detectionTickers[jobType] = ticker
go d.detectionLoop(jobType, ticker)
}
}
// Stop stops the dispatcher
func (d *Dispatcher) Stop() {
close(d.stopChan)
d.mu.Lock()
defer d.mu.Unlock()
for _, ticker := range d.detectionTickers {
ticker.Stop()
}
}
// detectionLoop periodically runs detection for a job type
func (d *Dispatcher) detectionLoop(jobType string, ticker *time.Ticker) {
for {
select {
case <-ticker.C:
d.runDetection(jobType)
case <-d.stopChan:
return
}
}
}
// runDetection runs detection for a job type and queues detected jobs
func (d *Dispatcher) runDetection(jobType string) {
// Get detector plugin
_, err := d.registry.GetDetectorForJobType(jobType)
if err != nil {
// No detector available, skip
return
}
// Get job type config
config, err := d.configMgr.GetConfig(jobType)
if err != nil || !config.Enabled {
// Config not found or disabled, skip
return
}
// TODO: Call detector plugin's DetectJobs method
// For now, this is a placeholder for the gRPC call
}
// QueueJob adds a job to the queue
func (d *Dispatcher) QueueJob(jobType string, req *plugin_pb.JobRequest, dedupKey string) (*Job, error) {
d.mu.RLock()
queue, exists := d.queues[jobType]
d.mu.RUnlock()
if !exists {
return nil, fmt.Errorf("job type %s not registered", jobType)
}
// Generate job ID
jobID := d.generateJobID()
return queue.AddJob(jobID, req, dedupKey)
}
// DispatchJob assigns a job to an available executor
func (d *Dispatcher) DispatchJob(jobType string) (*Job, *ConnectedPlugin, error) {
d.mu.RLock()
queue, exists := d.queues[jobType]
d.mu.RUnlock()
if !exists {
return nil, nil, fmt.Errorf("job type %s not registered", jobType)
}
// Get next pending job
job, err := queue.GetPendingJob()
if err != nil {
return nil, nil, err
}
// Get available executor
executor, err := d.registry.GetExecutorForJobType(jobType)
if err != nil {
// Put job back in queue
queue.mu.Lock()
queue.pendingQueue = append(queue.pendingQueue, job)
queue.mu.Unlock()
return nil, nil, err
}
job.ExecutorID = executor.ID
return job, executor, nil
}
// GetJob retrieves a job by ID across all queues
func (d *Dispatcher) GetJob(jobID string) (*Job, error) {
d.mu.RLock()
defer d.mu.RUnlock()
for _, queue := range d.queues {
job, err := queue.GetJob(jobID)
if err == nil {
return job, nil
}
}
return nil, fmt.Errorf("job %s not found", jobID)
}
// ListJobs returns jobs of a specific type and status
func (d *Dispatcher) ListJobs(jobType string, status JobState) []*Job {
d.mu.RLock()
queue, exists := d.queues[jobType]
d.mu.RUnlock()
if !exists {
return nil
}
return queue.ListJobs(status, 0)
}
// ListAllJobs returns all jobs across all queues
func (d *Dispatcher) ListAllJobs() []*Job {
d.mu.RLock()
defer d.mu.RUnlock()
var jobs []*Job
for _, queue := range d.queues {
jobs = append(jobs, queue.ListAllJobs()...)
}
return jobs
}
// UpdateJobProgress updates job progress
func (d *Dispatcher) UpdateJobProgress(jobID string, percent int32) error {
job, err := d.GetJob(jobID)
if err != nil {
return err
}
job.SetProgress(percent)
return nil
}
// CompleteJob marks a job as complete
func (d *Dispatcher) CompleteJob(jobID string, completed *plugin_pb.JobCompleted) error {
job, err := d.GetJob(jobID)
if err != nil {
return err
}
d.mu.RLock()
queue, exists := d.queues[job.Type]
d.mu.RUnlock()
if !exists {
return fmt.Errorf("job type %s not found", job.Type)
}
d.removeExecutionStream(jobID)
return queue.CompleteJob(jobID, completed)
}
// FailJob marks a job as failed
func (d *Dispatcher) FailJob(jobID string, failed *plugin_pb.JobFailed) error {
job, err := d.GetJob(jobID)
if err != nil {
return err
}
d.mu.RLock()
queue, exists := d.queues[job.Type]
d.mu.RUnlock()
if !exists {
return fmt.Errorf("job type %s not found", job.Type)
}
d.removeExecutionStream(jobID)
_, err = queue.FailJob(jobID, failed)
return err
}
// CancelJob cancels a job
func (d *Dispatcher) CancelJob(jobID string) error {
job, err := d.GetJob(jobID)
if err != nil {
return err
}
d.mu.RLock()
queue, exists := d.queues[job.Type]
d.mu.RUnlock()
if !exists {
return fmt.Errorf("job type %s not found", job.Type)
}
d.removeExecutionStream(jobID)
return queue.CancelJob(jobID)
}
// RetryJob retries a failed job
func (d *Dispatcher) RetryJob(jobID string) error {
job, err := d.GetJob(jobID)
if err != nil {
return err
}
state := job.GetState()
if state != JobStateFailed {
return fmt.Errorf("can only retry failed jobs, job is in %s state", state.String())
}
d.mu.RLock()
queue, exists := d.queues[job.Type]
d.mu.RUnlock()
if !exists {
return fmt.Errorf("job type %s not found", job.Type)
}
job.SetState(JobStatePending)
job.Retries++
queue.mu.Lock()
queue.pendingQueue = append(queue.pendingQueue, job)
queue.sortPendingQueue()
queue.mu.Unlock()
return nil
}
// generateJobID generates a unique job ID
func (d *Dispatcher) generateJobID() string {
d.jobCounterMu.Lock()
defer d.jobCounterMu.Unlock()
d.jobCounter++
return fmt.Sprintf("job-%d-%d", time.Now().Unix(), d.jobCounter)
}
// RegisterExecutionStream registers a new execution stream
func (d *Dispatcher) RegisterExecutionStream(jobID, pluginID string, ctx context.Context) *ExecutionStream {
d.streamsMu.Lock()
defer d.streamsMu.Unlock()
cancelCtx, cancel := context.WithCancel(ctx)
stream := &ExecutionStream{
JobID: jobID,
PluginID: pluginID,
Context: cancelCtx,
Cancel: cancel,
StartTime: time.Now(),
LastUpdate: time.Now(),
}
d.executionStreams[jobID] = stream
return stream
}
// removeExecutionStream removes an execution stream
func (d *Dispatcher) removeExecutionStream(jobID string) {
d.streamsMu.Lock()
defer d.streamsMu.Unlock()
stream, exists := d.executionStreams[jobID]
if exists {
stream.Cancel()
delete(d.executionStreams, jobID)
}
}
// GetStats returns dispatcher statistics
func (d *Dispatcher) GetStats() map[string]interface{} {
d.mu.RLock()
defer d.mu.RUnlock()
queueStats := make(map[string]interface{})
for jobType, queue := range d.queues {
queueStats[jobType] = queue.GetStats()
}
d.streamsMu.RLock()
activeStreams := len(d.executionStreams)
d.streamsMu.RUnlock()
return map[string]interface{}{
"queue_stats": queueStats,
"active_streams": activeStreams,
"registered_types": len(d.queues),
}
}
+461
View File
@@ -0,0 +1,461 @@
package plugin
import (
"context"
"fmt"
"io"
"sync"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
"google.golang.org/protobuf/types/known/emptypb"
"google.golang.org/grpc"
)
// GrpcServer implements the plugin gRPC service
type GrpcServer struct {
plugin_pb.UnimplementedPluginServiceServer
plugin_pb.UnimplementedAdminQueryServiceServer
plugin_pb.UnimplementedAdminCommandServiceServer
registry *Registry
dispatcher *Dispatcher
configMgr *ConfigManager
// Active connections
connections map[string]*grpc.ClientStream
connMu sync.RWMutex
}
// NewGrpcServer creates a new gRPC server
func NewGrpcServer(registry *Registry, dispatcher *Dispatcher, configMgr *ConfigManager) *GrpcServer {
return &GrpcServer{
registry: registry,
dispatcher: dispatcher,
configMgr: configMgr,
connections: make(map[string]*grpc.ClientStream),
}
}
// ============================================================================
// PluginService Implementation
// ============================================================================
// Connect handles bidirectional streaming for plugin registration and lifecycle
func (s *GrpcServer) Connect(stream grpc.BidiStreamingServer[plugin_pb.PluginMessage, plugin_pb.AdminMessage]) error {
var pluginID string
for {
msg, err := stream.Recv()
if err == io.EOF {
// Stream closed normally
if pluginID != "" {
s.registry.UnregisterPlugin(pluginID)
}
return nil
}
if err != nil {
return err
}
// Handle different message types
switch content := msg.Content.(type) {
case *plugin_pb.PluginMessage_Register:
register := content.Register
pluginID = register.PluginId
_, _ = s.registry.RegisterPlugin(pluginID, register)
if err != nil {
return err
}
case *plugin_pb.PluginMessage_Heartbeat:
heartbeat := content.Heartbeat
if err := s.registry.UpdateHeartbeat(heartbeat.PluginId, heartbeat); err != nil {
return err
}
case *plugin_pb.PluginMessage_StatusUpdate:
// Handle job status updates
statusUpdate := content.StatusUpdate
if err := s.handleStatusUpdate(statusUpdate); err != nil {
return err
}
}
// Send acknowledgment
adminMsg := &plugin_pb.AdminMessage{
Content: &plugin_pb.AdminMessage_AdminCommand{
AdminCommand: &plugin_pb.AdminCommand{
CommandType: plugin_pb.AdminCommand_RELOAD_CONFIG,
},
},
}
if err := stream.Send(adminMsg); err != nil {
if pluginID != "" {
s.registry.UnregisterPlugin(pluginID)
}
return err
}
}
}
// ExecuteJob handles bidirectional streaming for job execution
func (s *GrpcServer) ExecuteJob(stream grpc.BidiStreamingServer[plugin_pb.JobExecutionMessage, plugin_pb.JobProgressMessage]) error {
var jobID string
for {
execMsg, err := stream.Recv()
if err == io.EOF {
// Stream closed normally
return nil
}
if err != nil {
return err
}
jobID = execMsg.JobId
// Handle different execution message types
switch content := execMsg.Content.(type) {
case *plugin_pb.JobExecutionMessage_JobStarted:
started := content.JobStarted
if err := s.handleJobStarted(jobID, started); err != nil {
return err
}
case *plugin_pb.JobExecutionMessage_Progress:
progress := content.Progress
if err := s.dispatcher.UpdateJobProgress(jobID, progress.ProgressPercent); err != nil {
return err
}
case *plugin_pb.JobExecutionMessage_Checkpoint:
checkpoint := content.Checkpoint
// Store checkpoint for resumption
if err := s.storeCheckpoint(jobID, checkpoint); err != nil {
return err
}
case *plugin_pb.JobExecutionMessage_JobCompleted:
completed := content.JobCompleted
if err := s.dispatcher.CompleteJob(jobID, completed); err != nil {
return err
}
case *plugin_pb.JobExecutionMessage_JobFailed:
failed := content.JobFailed
if err := s.dispatcher.FailJob(jobID, failed); err != nil {
return err
}
case *plugin_pb.JobExecutionMessage_LogEntry:
logEntry := content.LogEntry
if err := s.storeLogEntry(jobID, logEntry); err != nil {
return err
}
}
// Send progress update
progressMsg := &plugin_pb.JobProgressMessage{
JobId: jobID,
Content: &plugin_pb.JobProgressMessage_Update{
Update: &plugin_pb.JobProgressUpdate{
JobId: jobID,
ProgressPercent: 0,
CurrentStep: "processing",
StatusMessage: "Processing job",
Timestamp: nil,
},
},
}
if err := stream.Send(progressMsg); err != nil {
return err
}
}
}
// ============================================================================
// AdminQueryService Implementation
// ============================================================================
// GetPluginStats returns current plugin system statistics
func (s *GrpcServer) GetPluginStats(ctx context.Context, req *emptypb.Empty) (*plugin_pb.PluginStats, error) {
plugins := s.registry.ListPlugins()
jobs := s.dispatcher.ListAllJobs()
var totalJobs, pendingJobs, runningJobs, completedJobs, failedJobs int32
for _, job := range jobs {
totalJobs++
switch job.GetState() {
case JobStatePending:
pendingJobs++
case JobStateRunning:
runningJobs++
case JobStateCompleted:
completedJobs++
case JobStateFailed:
failedJobs++
}
}
stats := &plugin_pb.PluginStats{
TotalPlugins: int32(len(plugins)),
ActivePlugins: int32(len(plugins)),
TotalJobs: totalJobs,
PendingJobs: pendingJobs,
RunningJobs: runningJobs,
CompletedJobs: completedJobs,
FailedJobs: failedJobs,
}
return stats, nil
}
// ListPlugins returns all connected plugins
func (s *GrpcServer) ListPlugins(ctx context.Context, req *emptypb.Empty) (*plugin_pb.PluginList, error) {
plugins := s.registry.ListPlugins()
pluginInfos := make([]*plugin_pb.PluginInfo, 0, len(plugins))
for _, p := range plugins {
info := &plugin_pb.PluginInfo{
PluginId: p.ID,
Name: p.Name,
Version: p.Version,
ProtocolVersion: p.ProtocolVersion,
Healthy: p.IsHealthy(),
Status: p.GetState().String(),
}
// Add capabilities
for _, cap := range p.Capabilities {
info.Capabilities = append(info.Capabilities, cap)
}
pluginInfos = append(pluginInfos, info)
}
return &plugin_pb.PluginList{Plugins: pluginInfos}, nil
}
// ListJobs returns jobs filtered by type and status
func (s *GrpcServer) ListJobs(ctx context.Context, req *plugin_pb.ListJobsRequest) (*plugin_pb.JobList, error) {
var state JobState
switch req.Status {
case plugin_pb.ListJobsRequest_PENDING:
state = JobStatePending
case plugin_pb.ListJobsRequest_RUNNING:
state = JobStateRunning
case plugin_pb.ListJobsRequest_COMPLETED:
state = JobStateCompleted
case plugin_pb.ListJobsRequest_FAILED:
state = JobStateFailed
default:
state = JobStatePending
}
jobs := s.dispatcher.ListJobs(req.JobType, state)
jobPbs := make([]*plugin_pb.Job, 0, len(jobs))
for _, job := range jobs {
jobPb := &plugin_pb.Job{
JobId: job.ID,
JobType: job.Type,
Description: job.Description,
Priority: job.Priority,
CreatedAt: nil,
UpdatedAt: nil,
ExecutorPluginId: job.ExecutorID,
ProgressPercent: job.GetProgress(),
Config: job.Config,
}
jobPbs = append(jobPbs, jobPb)
}
return &plugin_pb.JobList{
Jobs: jobPbs,
TotalCount: int32(len(jobPbs)),
ReturnedCount: int32(len(jobPbs)),
}, nil
}
// GetJob returns detailed job information
func (s *GrpcServer) GetJob(ctx context.Context, req *plugin_pb.GetJobRequest) (*plugin_pb.JobDetail, error) {
job, err := s.dispatcher.GetJob(req.JobId)
if err != nil {
return nil, fmt.Errorf("job not found: %w", err)
}
jobPb := &plugin_pb.Job{
JobId: job.ID,
JobType: job.Type,
Description: job.Description,
Priority: job.Priority,
ExecutorPluginId: job.ExecutorID,
ProgressPercent: job.GetProgress(),
Config: job.Config,
}
detail := &plugin_pb.JobDetail{
Job: jobPb,
LogEntries: job.LogEntries,
CompletionInfo: job.CompletedInfo,
FailureInfo: job.FailedInfo,
CheckpointIds: job.CheckpointIDs,
}
return detail, nil
}
// GetJobExecutionLog returns job execution logs
func (s *GrpcServer) GetJobExecutionLog(req *plugin_pb.GetJobExecutionLogRequest, stream grpc.ServerStreamingServer[plugin_pb.ExecutionLogEntry]) error {
job, err := s.dispatcher.GetJob(req.JobId)
if err != nil {
return fmt.Errorf("job not found: %w", err)
}
// Send log entries
for _, entry := range job.LogEntries {
if err := stream.Send(entry); err != nil {
return err
}
}
return nil
}
// ============================================================================
// AdminCommandService Implementation
// ============================================================================
// UpdateJobTypeConfig updates configuration for a job type
func (s *GrpcServer) UpdateJobTypeConfig(ctx context.Context, req *plugin_pb.UpdateJobTypeConfigRequest) (*plugin_pb.JobTypeConfig, error) {
config := req.Config
_, err := s.configMgr.UpdateConfig(
config.JobType,
config.Enabled,
config.AdminConfig,
config.WorkerConfig,
"admin",
)
if err != nil {
return nil, fmt.Errorf("failed to update config: %w", err)
}
return config, nil
}
// CreateJob creates a new job
func (s *GrpcServer) CreateJob(ctx context.Context, req *plugin_pb.CreateJobRequest) (*plugin_pb.Job, error) {
jobReq := &plugin_pb.JobRequest{
JobType: req.JobType,
Description: req.Description,
Priority: req.Priority,
Config: req.Config,
Metadata: req.Metadata,
RequestSource: "admin",
}
job, err := s.dispatcher.QueueJob(req.JobType, jobReq, "")
if err != nil {
return nil, fmt.Errorf("failed to create job: %w", err)
}
return &plugin_pb.Job{
JobId: job.ID,
JobType: job.Type,
Description: job.Description,
Priority: job.Priority,
ProgressPercent: job.GetProgress(),
Config: job.Config,
}, nil
}
// CancelJob cancels a running job
func (s *GrpcServer) CancelJob(ctx context.Context, req *plugin_pb.CancelJobRequest) (*emptypb.Empty, error) {
if err := s.dispatcher.CancelJob(req.JobId); err != nil {
return nil, fmt.Errorf("failed to cancel job: %w", err)
}
return &emptypb.Empty{}, nil
}
// RetryJob retries a failed job
func (s *GrpcServer) RetryJob(ctx context.Context, req *plugin_pb.RetryJobRequest) (*plugin_pb.Job, error) {
if err := s.dispatcher.RetryJob(req.JobId); err != nil {
return nil, fmt.Errorf("failed to retry job: %w", err)
}
job, err := s.dispatcher.GetJob(req.JobId)
if err != nil {
return nil, fmt.Errorf("job not found: %w", err)
}
return &plugin_pb.Job{
JobId: job.ID,
JobType: job.Type,
Description: job.Description,
Priority: job.Priority,
ProgressPercent: job.GetProgress(),
Config: job.Config,
}, nil
}
// ============================================================================
// Helper Methods
// ============================================================================
// handleStatusUpdate handles execution status updates
func (s *GrpcServer) handleStatusUpdate(update *plugin_pb.ExecutionStatusUpdate) error {
if err := s.dispatcher.UpdateJobProgress(update.JobId, update.ProgressPercent); err != nil {
return err
}
return nil
}
// handleJobStarted handles job start messages
func (s *GrpcServer) handleJobStarted(jobID string, started *plugin_pb.JobStarted) error {
job, err := s.dispatcher.GetJob(jobID)
if err != nil {
return err
}
job.SetState(JobStateRunning)
return nil
}
// storeCheckpoint stores a job checkpoint
func (s *GrpcServer) storeCheckpoint(jobID string, checkpoint *plugin_pb.JobCheckpoint) error {
job, err := s.dispatcher.GetJob(jobID)
if err != nil {
return err
}
job.mu.Lock()
defer job.mu.Unlock()
job.CheckpointIDs = append(job.CheckpointIDs, checkpoint.CheckpointId)
return nil
}
// storeLogEntry stores a log entry for a job
func (s *GrpcServer) storeLogEntry(jobID string, logEntry *plugin_pb.ExecutionLog) error {
job, err := s.dispatcher.GetJob(jobID)
if err != nil {
return err
}
job.mu.Lock()
defer job.mu.Unlock()
job.LogEntries = append(job.LogEntries, &plugin_pb.ExecutionLogEntry{
Timestamp: logEntry.Timestamp,
Level: logEntry.Level,
Message: logEntry.Message,
Context: logEntry.Context,
})
return nil
}
+346
View File
@@ -0,0 +1,346 @@
package plugin
import (
"fmt"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
)
// JobQueue manages jobs for a specific job type with priority and deduplication
type JobQueue struct {
jobType string
// Active jobs by ID
jobs map[string]*Job
// Priority queue
pendingQueue []*Job
// Completed jobs history (keeps last 1000)
completedHistory []*Job
maxHistorySize int
// Deduplication tracking
seenKeys map[string]bool
mu sync.RWMutex
// Retry configuration
maxRetries int32
retryBackoff time.Duration
retryBackoffMax time.Duration
}
// NewJobQueue creates a new job queue for a job type
func NewJobQueue(jobType string) *JobQueue {
return &JobQueue{
jobType: jobType,
jobs: make(map[string]*Job),
pendingQueue: make([]*Job, 0),
completedHistory: make([]*Job, 0),
maxHistorySize: 1000,
seenKeys: make(map[string]bool),
maxRetries: 3,
retryBackoff: time.Second,
retryBackoffMax: 5 * time.Minute,
}
}
// AddJob adds a new job to the queue with deduplication
func (q *JobQueue) AddJob(jobID string, req *plugin_pb.JobRequest, dedupKey string) (*Job, error) {
q.mu.Lock()
defer q.mu.Unlock()
// Check deduplication
if dedupKey != "" && q.seenKeys[dedupKey] {
return nil, fmt.Errorf("job with key %s already queued", dedupKey)
}
if _, exists := q.jobs[jobID]; exists {
return nil, fmt.Errorf("job %s already exists", jobID)
}
job := &Job{
ID: jobID,
Type: req.JobType,
Description: req.Description,
Priority: req.Priority,
State: JobStatePending,
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
Config: req.Config,
Metadata: req.Metadata,
Retries: 0,
}
q.jobs[jobID] = job
q.pendingQueue = append(q.pendingQueue, job)
if dedupKey != "" {
q.seenKeys[dedupKey] = true
}
// Sort by priority (higher priority first)
q.sortPendingQueue()
return job, nil
}
// GetPendingJob returns the next job to execute
func (q *JobQueue) GetPendingJob() (*Job, error) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.pendingQueue) == 0 {
return nil, fmt.Errorf("no pending jobs")
}
job := q.pendingQueue[0]
q.pendingQueue = q.pendingQueue[1:]
job.SetState(JobStateRunning)
return job, nil
}
// GetJob retrieves a job by ID
func (q *JobQueue) GetJob(jobID string) (*Job, error) {
q.mu.RLock()
defer q.mu.RUnlock()
job, exists := q.jobs[jobID]
if !exists {
return nil, fmt.Errorf("job %s not found", jobID)
}
return job, nil
}
// CompleteJob marks a job as completed
func (q *JobQueue) CompleteJob(jobID string, completed *plugin_pb.JobCompleted) error {
q.mu.Lock()
defer q.mu.Unlock()
job, exists := q.jobs[jobID]
if !exists {
return fmt.Errorf("job %s not found", jobID)
}
job.SetState(JobStateCompleted)
job.CompletedInfo = completed
q.addToHistory(job)
return nil
}
// FailJob marks a job as failed and determines if it should be retried
func (q *JobQueue) FailJob(jobID string, failed *plugin_pb.JobFailed) (*Job, error) {
q.mu.Lock()
defer q.mu.Unlock()
job, exists := q.jobs[jobID]
if !exists {
return nil, fmt.Errorf("job %s not found", jobID)
}
job.FailedInfo = failed
// Check if we should retry
if failed.Retryable && job.Retries < q.maxRetries {
job.Retries++
job.LastRetryTime = time.Now()
job.SetState(JobStatePending)
// Add back to queue with exponential backoff for retry
q.pendingQueue = append(q.pendingQueue, job)
q.sortPendingQueue()
return job, nil
}
// No retry, mark as failed
job.SetState(JobStateFailed)
q.addToHistory(job)
return job, nil
}
// PauseJob pauses a running job
func (q *JobQueue) PauseJob(jobID string) error {
q.mu.Lock()
defer q.mu.Unlock()
job, exists := q.jobs[jobID]
if !exists {
return fmt.Errorf("job %s not found", jobID)
}
if job.GetState() != JobStateRunning {
return fmt.Errorf("job %s is not running", jobID)
}
job.SetState(JobStatePaused)
return nil
}
// ResumeJob resumes a paused job
func (q *JobQueue) ResumeJob(jobID string) error {
q.mu.Lock()
defer q.mu.Unlock()
job, exists := q.jobs[jobID]
if !exists {
return fmt.Errorf("job %s not found", jobID)
}
if job.GetState() != JobStatePaused {
return fmt.Errorf("job %s is not paused", jobID)
}
job.SetState(JobStateRunning)
return nil
}
// CancelJob cancels a job
func (q *JobQueue) CancelJob(jobID string) error {
q.mu.Lock()
defer q.mu.Unlock()
job, exists := q.jobs[jobID]
if !exists {
return fmt.Errorf("job %s not found", jobID)
}
state := job.GetState()
if state != JobStatePending && state != JobStateRunning && state != JobStatePaused {
return fmt.Errorf("cannot cancel job in %s state", state.String())
}
job.SetState(JobStateCancelled)
// Remove from pending queue if present
for i, j := range q.pendingQueue {
if j.ID == jobID {
q.pendingQueue = append(q.pendingQueue[:i], q.pendingQueue[i+1:]...)
break
}
}
q.addToHistory(job)
return nil
}
// ListJobs returns jobs filtered by status
func (q *JobQueue) ListJobs(status JobState, limit int) []*Job {
q.mu.RLock()
defer q.mu.RUnlock()
var jobs []*Job
// Include pending jobs from queue
for _, job := range q.pendingQueue {
if job.GetState() == status {
jobs = append(jobs, job)
}
}
// Include all active jobs
for _, job := range q.jobs {
if job.GetState() == status {
jobs = append(jobs, job)
}
}
if limit > 0 && len(jobs) > limit {
jobs = jobs[:limit]
}
return jobs
}
// ListAllJobs returns all jobs
func (q *JobQueue) ListAllJobs() []*Job {
q.mu.RLock()
defer q.mu.RUnlock()
jobs := make([]*Job, 0, len(q.jobs))
for _, job := range q.jobs {
jobs = append(jobs, job)
}
// Add completed history
for _, job := range q.completedHistory {
jobs = append(jobs, job)
}
return jobs
}
// sortPendingQueue sorts the pending queue by priority (higher first)
func (q *JobQueue) sortPendingQueue() {
for i := 0; i < len(q.pendingQueue)-1; i++ {
for j := i + 1; j < len(q.pendingQueue); j++ {
if q.pendingQueue[j].Priority > q.pendingQueue[i].Priority {
q.pendingQueue[i], q.pendingQueue[j] = q.pendingQueue[j], q.pendingQueue[i]
}
}
}
}
// addToHistory adds a completed job to history
func (q *JobQueue) addToHistory(job *Job) {
q.completedHistory = append(q.completedHistory, job)
// Keep only last N entries
if len(q.completedHistory) > q.maxHistorySize {
q.completedHistory = q.completedHistory[len(q.completedHistory)-q.maxHistorySize:]
}
// Remove from active jobs
delete(q.jobs, job.ID)
}
// GetPendingCount returns number of pending jobs
func (q *JobQueue) GetPendingCount() int {
q.mu.RLock()
defer q.mu.RUnlock()
return len(q.pendingQueue)
}
// GetStats returns queue statistics
func (q *JobQueue) GetStats() map[string]interface{} {
q.mu.RLock()
defer q.mu.RUnlock()
var running, completed, failed, cancelled int
for _, job := range q.jobs {
switch job.GetState() {
case JobStateRunning:
running++
case JobStateCompleted:
completed++
case JobStateFailed:
failed++
case JobStateCancelled:
cancelled++
}
}
completed += len(q.completedHistory)
return map[string]interface{}{
"job_type": q.jobType,
"pending": len(q.pendingQueue),
"running": running,
"completed": completed,
"failed": failed,
"cancelled": cancelled,
"total": len(q.jobs) + len(q.completedHistory),
"history_size": len(q.completedHistory),
}
}
+401
View File
@@ -0,0 +1,401 @@
package plugin
import (
"context"
"fmt"
"net"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
"google.golang.org/grpc"
)
// Manager is the main plugin system manager
type Manager struct {
registry *Registry
dispatcher *Dispatcher
configMgr *ConfigManager
grpcServer *GrpcServer
// gRPC server
grpcListener net.Listener
grpc *grpc.Server
// Configuration
listenAddr string
dataDir string
healthCheckInterval time.Duration
// State
mu sync.RWMutex
running bool
stopChan chan struct{}
}
// NewManager creates a new plugin manager
func NewManager(listenAddr, dataDir string) (*Manager, error) {
// Create configuration manager
configMgr, err := NewConfigManager(dataDir)
if err != nil {
return nil, fmt.Errorf("failed to create config manager: %w", err)
}
// Create registry
registry := NewRegistry()
// Create dispatcher
dispatcher := NewDispatcher(registry, configMgr)
// Create gRPC server implementation
grpcServer := NewGrpcServer(registry, dispatcher, configMgr)
return &Manager{
registry: registry,
dispatcher: dispatcher,
configMgr: configMgr,
grpcServer: grpcServer,
listenAddr: listenAddr,
dataDir: dataDir,
healthCheckInterval: 30 * time.Second,
stopChan: make(chan struct{}),
}, nil
}
// Start starts the plugin manager and gRPC server
func (m *Manager) Start() error {
m.mu.Lock()
defer m.mu.Unlock()
if m.running {
return fmt.Errorf("plugin manager already running")
}
// Start registry health checks
m.registry.Start()
// Start dispatcher
m.dispatcher.Start()
// Create gRPC server
m.grpc = grpc.NewServer()
// Register services
plugin_pb.RegisterPluginServiceServer(m.grpc, m.grpcServer)
plugin_pb.RegisterAdminQueryServiceServer(m.grpc, m.grpcServer)
plugin_pb.RegisterAdminCommandServiceServer(m.grpc, m.grpcServer)
// Start listening
listener, err := net.Listen("tcp", m.listenAddr)
if err != nil {
return fmt.Errorf("failed to listen on %s: %w", m.listenAddr, err)
}
m.grpcListener = listener
// Start gRPC server in background
go func() {
if err := m.grpc.Serve(listener); err != nil && err != grpc.ErrServerStopped {
fmt.Printf("gRPC server error: %v\n", err)
}
}()
m.running = true
fmt.Printf("Plugin manager started on %s\n", m.listenAddr)
return nil
}
// Stop stops the plugin manager and gRPC server
func (m *Manager) Stop() error {
m.mu.Lock()
defer m.mu.Unlock()
if !m.running {
return fmt.Errorf("plugin manager not running")
}
close(m.stopChan)
// Stop dispatcher
m.dispatcher.Stop()
// Stop registry
m.registry.Stop()
// Stop gRPC server
if m.grpc != nil {
m.grpc.GracefulStop()
}
m.running = false
return nil
}
// IsRunning returns whether the manager is running
func (m *Manager) IsRunning() bool {
m.mu.RLock()
defer m.mu.RUnlock()
return m.running
}
// RegisterJobType registers a job type for detection and execution
func (m *Manager) RegisterJobType(jobType string, detectionInterval time.Duration) error {
m.mu.RLock()
running := m.running
m.mu.RUnlock()
if !running {
return fmt.Errorf("plugin manager not running")
}
// Register job type
if err := m.dispatcher.RegisterJobType(jobType, detectionInterval); err != nil {
return err
}
_, err := m.configMgr.LoadConfig(jobType)
if err != nil {
return err
}
return nil
}
// QueueJob queues a new job for execution
func (m *Manager) QueueJob(jobType string, description string, priority int64, metadata map[string]string) (*Job, error) {
m.mu.RLock()
running := m.running
m.mu.RUnlock()
if !running {
return nil, fmt.Errorf("plugin manager not running")
}
jobReq := &plugin_pb.JobRequest{
JobType: jobType,
Description: description,
Priority: priority,
Metadata: metadata,
RequestSource: "admin",
}
return m.dispatcher.QueueJob(jobType, jobReq, "")
}
// DispatchJob dispatches the next available job
func (m *Manager) DispatchJob(jobType string) (*Job, *ConnectedPlugin, error) {
m.mu.RLock()
running := m.running
m.mu.RUnlock()
if !running {
return nil, nil, fmt.Errorf("plugin manager not running")
}
return m.dispatcher.DispatchJob(jobType)
}
// GetJob retrieves a job by ID
func (m *Manager) GetJob(jobID string) (*Job, error) {
return m.dispatcher.GetJob(jobID)
}
// ListJobs returns jobs of a specific type and status
func (m *Manager) ListJobs(jobType string, status JobState) []*Job {
return m.dispatcher.ListJobs(jobType, status)
}
// ListAllJobs returns all jobs
func (m *Manager) ListAllJobs() []*Job {
return m.dispatcher.ListAllJobs()
}
// ListPlugins returns all connected plugins
func (m *Manager) ListPlugins() []*ConnectedPlugin {
return m.registry.ListPlugins()
}
// GetPlugin retrieves a plugin by ID
func (m *Manager) GetPlugin(pluginID string) (*ConnectedPlugin, error) {
return m.registry.GetPlugin(pluginID)
}
// ListPluginsByCapability returns plugins that support a job type
func (m *Manager) ListPluginsByCapability(jobType string) []*ConnectedPlugin {
return m.registry.ListPluginsByCapability(jobType)
}
// UpdateJobTypeConfig updates configuration for a job type
func (m *Manager) UpdateJobTypeConfig(jobType string, enabled bool, config []*plugin_pb.ConfigFieldValue) (*JobTypeConfig, error) {
return m.configMgr.UpdateConfig(jobType, enabled, config, nil, "admin")
}
// GetJobTypeConfig retrieves configuration for a job type
func (m *Manager) GetJobTypeConfig(jobType string) (*JobTypeConfig, error) {
return m.configMgr.GetConfig(jobType)
}
// CancelJob cancels a job
func (m *Manager) CancelJob(jobID string) error {
return m.dispatcher.CancelJob(jobID)
}
// RetryJob retries a failed job
func (m *Manager) RetryJob(jobID string) error {
return m.dispatcher.RetryJob(jobID)
}
// GetStats returns system statistics
func (m *Manager) GetStats() map[string]interface{} {
m.mu.RLock()
running := m.running
m.mu.RUnlock()
stats := map[string]interface{}{
"running": running,
}
if running {
stats["registry"] = m.registry.GetStats()
stats["dispatcher"] = m.dispatcher.GetStats()
}
return stats
}
// WaitForPluginConnection waits for a plugin to connect with timeout
func (m *Manager) WaitForPluginConnection(timeout time.Duration) error {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
plugins := m.registry.ListPlugins()
if len(plugins) > 0 {
return nil
}
time.Sleep(100 * time.Millisecond)
}
return fmt.Errorf("no plugin connected within %v", timeout)
}
// CreateTestJob creates a job for testing
func (m *Manager) CreateTestJob(jobType string, description string) (*Job, error) {
return m.QueueJob(jobType, description, 100, map[string]string{
"test": "true",
})
}
// GetHealthStatus returns the health status of the plugin system
func (m *Manager) GetHealthStatus() map[string]interface{} {
m.mu.RLock()
running := m.running
m.mu.RUnlock()
status := map[string]interface{}{
"running": running,
}
if !running {
return status
}
plugins := m.registry.ListPlugins()
var healthyCount int
for _, p := range plugins {
if p.IsHealthy() {
healthyCount++
}
}
jobs := m.dispatcher.ListAllJobs()
var runningJobs, pendingJobs, completedJobs, failedJobs int
for _, job := range jobs {
switch job.GetState() {
case JobStateRunning:
runningJobs++
case JobStatePending:
pendingJobs++
case JobStateCompleted:
completedJobs++
case JobStateFailed:
failedJobs++
}
}
status["plugins"] = map[string]interface{}{
"total": len(plugins),
"healthy": healthyCount,
}
status["jobs"] = map[string]interface{}{
"running": runningJobs,
"pending": pendingJobs,
"completed": completedJobs,
"failed": failedJobs,
"total": len(jobs),
}
return status
}
// LoadAllConfigs loads all configurations from disk
func (m *Manager) LoadAllConfigs() error {
return m.configMgr.LoadAllConfigs()
}
// GetDataDir returns the data directory
func (m *Manager) GetDataDir() string {
return m.configMgr.GetDataDir()
}
// QueryPluginService queries a plugin service for capabilities
func (m *Manager) QueryPluginService(ctx context.Context, pluginID string, addr string) (*plugin_pb.PluginRegister, error) {
// This would be implemented to connect to a plugin and query its capabilities
// For now, return a placeholder error
return nil, fmt.Errorf("not yet implemented")
}
// ExecuteDetection triggers detection for a job type
func (m *Manager) ExecuteDetection(ctx context.Context, jobType string) ([]*DetectedJobInfo, error) {
_, err := m.registry.GetDetectorForJobType(jobType)
if err != nil {
return nil, fmt.Errorf("no detector available for job type %s: %w", jobType, err)
}
// TODO: Call detector plugin's DetectJobs method
// This would require connecting to the plugin and calling its service
return nil, fmt.Errorf("detection not yet implemented")
}
// DetectedJobInfo holds information about a detected job
type DetectedJobInfo struct {
Key string
JobType string
Priority int64
Metadata map[string]string
Timestamp time.Time
}
// GetListenAddr returns the listen address
func (m *Manager) GetListenAddr() string {
m.mu.RLock()
defer m.mu.RUnlock()
return m.listenAddr
}
// GetContext returns a context for the manager
func (m *Manager) GetContext() context.Context {
return context.Background()
}
// Close closes the manager and releases resources
func (m *Manager) Close() error {
if m.IsRunning() {
return m.Stop()
}
return nil
}
+333
View File
@@ -0,0 +1,333 @@
package plugin
import (
"fmt"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
)
// Registry manages all connected plugins and their capabilities
type Registry struct {
plugins map[string]*ConnectedPlugin
// Index by job type and capability (detector/executor)
detectorsByJobType map[string][]*ConnectedPlugin
executorsByJobType map[string][]*ConnectedPlugin
mu sync.RWMutex
// Configuration
healthCheckInterval time.Duration
healthCheckTimeout time.Duration
healthCheckTicker *time.Ticker
stopChan chan struct{}
}
// NewRegistry creates a new plugin registry
func NewRegistry() *Registry {
return &Registry{
plugins: make(map[string]*ConnectedPlugin),
detectorsByJobType: make(map[string][]*ConnectedPlugin),
executorsByJobType: make(map[string][]*ConnectedPlugin),
healthCheckInterval: 30 * time.Second,
healthCheckTimeout: 90 * time.Second,
stopChan: make(chan struct{}),
}
}
// RegisterPlugin adds a new plugin to the registry
func (r *Registry) RegisterPlugin(pluginID string, register *plugin_pb.PluginRegister) (*ConnectedPlugin, error) {
r.mu.Lock()
defer r.mu.Unlock()
if _, exists := r.plugins[pluginID]; exists {
return nil, fmt.Errorf("plugin %s already registered", pluginID)
}
plugin := &ConnectedPlugin{
ID: pluginID,
Name: register.Name,
Version: register.Version,
ProtocolVersion: register.ProtocolVersion,
ConnectedAt: time.Now(),
LastHeartbeat: time.Now(),
State: PluginStateConnected,
Healthy: true,
Capabilities: make(map[string]*plugin_pb.JobTypeCapability),
StreamConnected: true,
}
// Index capabilities
for _, cap := range register.Capabilities {
plugin.Capabilities[cap.JobType] = cap
if cap.CanDetect {
r.detectorsByJobType[cap.JobType] = append(r.detectorsByJobType[cap.JobType], plugin)
}
if cap.CanExecute {
r.executorsByJobType[cap.JobType] = append(r.executorsByJobType[cap.JobType], plugin)
}
}
r.plugins[pluginID] = plugin
return plugin, nil
}
// UnregisterPlugin removes a plugin from the registry
func (r *Registry) UnregisterPlugin(pluginID string) error {
r.mu.Lock()
defer r.mu.Unlock()
plugin, exists := r.plugins[pluginID]
if !exists {
return fmt.Errorf("plugin %s not found", pluginID)
}
// Remove from capability indexes
for jobType := range plugin.Capabilities {
r.removePluginFromDetectors(jobType, plugin)
r.removePluginFromExecutors(jobType, plugin)
}
delete(r.plugins, pluginID)
plugin.SetState(PluginStateDisconnected)
return nil
}
// removePluginFromDetectors removes plugin from detector index for a job type
func (r *Registry) removePluginFromDetectors(jobType string, plugin *ConnectedPlugin) {
detectors := r.detectorsByJobType[jobType]
for i, p := range detectors {
if p.ID == plugin.ID {
r.detectorsByJobType[jobType] = append(detectors[:i], detectors[i+1:]...)
break
}
}
if len(r.detectorsByJobType[jobType]) == 0 {
delete(r.detectorsByJobType, jobType)
}
}
// removePluginFromExecutors removes plugin from executor index for a job type
func (r *Registry) removePluginFromExecutors(jobType string, plugin *ConnectedPlugin) {
executors := r.executorsByJobType[jobType]
for i, p := range executors {
if p.ID == plugin.ID {
r.executorsByJobType[jobType] = append(executors[:i], executors[i+1:]...)
break
}
}
if len(r.executorsByJobType[jobType]) == 0 {
delete(r.executorsByJobType, jobType)
}
}
// GetPlugin retrieves a plugin by ID
func (r *Registry) GetPlugin(pluginID string) (*ConnectedPlugin, error) {
r.mu.RLock()
defer r.mu.RUnlock()
plugin, exists := r.plugins[pluginID]
if !exists {
return nil, fmt.Errorf("plugin %s not found", pluginID)
}
return plugin, nil
}
// GetDetectorForJobType returns an available detector plugin for a job type
func (r *Registry) GetDetectorForJobType(jobType string) (*ConnectedPlugin, error) {
r.mu.RLock()
defer r.mu.RUnlock()
detectors, exists := r.detectorsByJobType[jobType]
if !exists || len(detectors) == 0 {
return nil, fmt.Errorf("no detector available for job type %s", jobType)
}
// Return the first healthy detector
for _, detector := range detectors {
if detector.IsHealthy() {
return detector, nil
}
}
return nil, fmt.Errorf("no healthy detector available for job type %s", jobType)
}
// GetExecutorForJobType returns an available executor plugin for a job type
func (r *Registry) GetExecutorForJobType(jobType string) (*ConnectedPlugin, error) {
r.mu.RLock()
defer r.mu.RUnlock()
executors, exists := r.executorsByJobType[jobType]
if !exists || len(executors) == 0 {
return nil, fmt.Errorf("no executor available for job type %s", jobType)
}
// Return the first healthy executor with lowest workload
var bestExecutor *ConnectedPlugin
var minLoad int32 = int32(^uint32(0) >> 1)
for _, executor := range executors {
if executor.IsHealthy() {
load := executor.PendingJobs + executor.RunningJobs
if load < minLoad {
minLoad = load
bestExecutor = executor
}
}
}
if bestExecutor == nil {
return nil, fmt.Errorf("no healthy executor available for job type %s", jobType)
}
return bestExecutor, nil
}
// ListPlugins returns all connected plugins
func (r *Registry) ListPlugins() []*ConnectedPlugin {
r.mu.RLock()
defer r.mu.RUnlock()
plugins := make([]*ConnectedPlugin, 0, len(r.plugins))
for _, plugin := range r.plugins {
plugins = append(plugins, plugin)
}
return plugins
}
// ListPluginsByCapability returns plugins that support a specific job type
func (r *Registry) ListPluginsByCapability(jobType string) []*ConnectedPlugin {
r.mu.RLock()
defer r.mu.RUnlock()
detectors := r.detectorsByJobType[jobType]
executors := r.executorsByJobType[jobType]
// Combine and deduplicate
pluginMap := make(map[string]*ConnectedPlugin)
for _, p := range detectors {
pluginMap[p.ID] = p
}
for _, p := range executors {
pluginMap[p.ID] = p
}
plugins := make([]*ConnectedPlugin, 0, len(pluginMap))
for _, p := range pluginMap {
plugins = append(plugins, p)
}
return plugins
}
// UpdateHeartbeat updates plugin heartbeat and health status
func (r *Registry) UpdateHeartbeat(pluginID string, heartbeat *plugin_pb.PluginHeartbeat) error {
r.mu.Lock()
plugin, exists := r.plugins[pluginID]
r.mu.Unlock()
if !exists {
return fmt.Errorf("plugin %s not found", pluginID)
}
plugin.UpdateHeartbeat(
heartbeat.CpuUsagePercent,
heartbeat.MemoryUsageMb,
heartbeat.PendingJobs,
int32(len(plugin.Capabilities)), // rough estimate
)
return nil
}
// Start starts the health check loop
func (r *Registry) Start() {
r.healthCheckTicker = time.NewTicker(r.healthCheckInterval)
go r.healthCheckLoop()
}
// Stop stops the health check loop
func (r *Registry) Stop() {
close(r.stopChan)
if r.healthCheckTicker != nil {
r.healthCheckTicker.Stop()
}
}
// healthCheckLoop periodically checks plugin health
func (r *Registry) healthCheckLoop() {
for {
select {
case <-r.healthCheckTicker.C:
r.checkHealth()
case <-r.stopChan:
return
}
}
}
// checkHealth checks the health of all plugins
func (r *Registry) checkHealth() {
r.mu.RLock()
plugins := make([]*ConnectedPlugin, 0, len(r.plugins))
for _, p := range r.plugins {
plugins = append(plugins, p)
}
r.mu.RUnlock()
now := time.Now()
for _, plugin := range plugins {
plugin.mu.RLock()
lastHeartbeat := plugin.LastHeartbeat
plugin.mu.RUnlock()
timeSinceHeartbeat := now.Sub(lastHeartbeat)
if timeSinceHeartbeat > r.healthCheckTimeout {
plugin.SetState(PluginStateUnhealthy)
plugin.mu.Lock()
plugin.Healthy = false
plugin.mu.Unlock()
} else if timeSinceHeartbeat > r.healthCheckTimeout/2 {
plugin.SetState(PluginStateUnhealthy)
} else {
plugin.SetState(PluginStateHealthy)
plugin.mu.Lock()
plugin.Healthy = true
plugin.mu.Unlock()
}
}
}
// GetStats returns registry statistics
func (r *Registry) GetStats() map[string]interface{} {
r.mu.RLock()
defer r.mu.RUnlock()
stats := map[string]interface{}{
"total_plugins": len(r.plugins),
"job_types": len(r.detectorsByJobType),
}
var healthyCount int
for _, p := range r.plugins {
if p.IsHealthy() {
healthyCount++
}
}
stats["healthy_plugins"] = healthyCount
jobTypeStats := make(map[string]map[string]interface{})
for jobType := range r.detectorsByJobType {
jobTypeStats[jobType] = map[string]interface{}{
"detectors": len(r.detectorsByJobType[jobType]),
"executors": len(r.executorsByJobType[jobType]),
}
}
stats["job_type_stats"] = jobTypeStats
return stats
}
+227
View File
@@ -0,0 +1,227 @@
package plugin
import (
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// JobState represents the current state of a job
type JobState int32
const (
JobStatePending JobState = 0
JobStateRunning JobState = 1
JobStateCompleted JobState = 2
JobStateFailed JobState = 3
JobStateCancelled JobState = 4
JobStatePaused JobState = 5
)
// String returns string representation of JobState
func (s JobState) String() string {
switch s {
case JobStatePending:
return "pending"
case JobStateRunning:
return "running"
case JobStateCompleted:
return "completed"
case JobStateFailed:
return "failed"
case JobStateCancelled:
return "cancelled"
case JobStatePaused:
return "paused"
default:
return "unknown"
}
}
// PluginState represents the current state of a plugin connection
type PluginState int32
const (
PluginStateConnecting PluginState = 0
PluginStateConnected PluginState = 1
PluginStateHealthy PluginState = 2
PluginStateUnhealthy PluginState = 3
PluginStateDisconnected PluginState = 4
)
// String returns string representation of PluginState
func (s PluginState) String() string {
switch s {
case PluginStateConnecting:
return "connecting"
case PluginStateConnected:
return "connected"
case PluginStateHealthy:
return "healthy"
case PluginStateUnhealthy:
return "unhealthy"
case PluginStateDisconnected:
return "disconnected"
default:
return "unknown"
}
}
// Job represents a work unit in the plugin system
type Job struct {
ID string
Type string
Description string
Priority int64
State JobState
CreatedAt time.Time
UpdatedAt time.Time
ExecutorID string
ProgressPercent int32
Config []*plugin_pb.ConfigFieldValue
Metadata map[string]string
// Execution tracking
Retries int32
LastRetryTime time.Time
CompletedInfo *plugin_pb.JobCompleted
FailedInfo *plugin_pb.JobFailed
CheckpointIDs []string
LogEntries []*plugin_pb.ExecutionLogEntry
// Lock for thread-safe access
mu sync.RWMutex
}
// ConnectedPlugin represents a connected plugin worker
type ConnectedPlugin struct {
ID string
Name string
Version string
ProtocolVersion string
ConnectedAt time.Time
LastHeartbeat time.Time
State PluginState
Healthy bool
// Capabilities indexed by job type
Capabilities map[string]*plugin_pb.JobTypeCapability
// Current workload
PendingJobs int32
RunningJobs int32
// Resource usage
CPUUsagePercent float32
MemoryUsageMB float32
// Communication stream
StreamConnected bool
// Lock for thread-safe access
mu sync.RWMutex
}
// JobTypeConfig holds configuration for a specific job type
type JobTypeConfig struct {
JobType string
Enabled bool
AdminConfig []*plugin_pb.ConfigFieldValue
WorkerConfig []*plugin_pb.ConfigFieldValue
CreatedAt time.Time
UpdatedAt time.Time
CreatedBy string
mu sync.RWMutex
}
// GetState safely gets the job state
func (j *Job) GetState() JobState {
j.mu.RLock()
defer j.mu.RUnlock()
return j.State
}
// SetState safely sets the job state
func (j *Job) SetState(state JobState) {
j.mu.Lock()
defer j.mu.Unlock()
j.State = state
j.UpdatedAt = time.Now()
}
// GetProgress safely gets progress
func (j *Job) GetProgress() int32 {
j.mu.RLock()
defer j.mu.RUnlock()
return j.ProgressPercent
}
// SetProgress safely sets progress
func (j *Job) SetProgress(percent int32) {
j.mu.Lock()
defer j.mu.Unlock()
j.ProgressPercent = percent
j.UpdatedAt = time.Now()
}
// GetState safely gets plugin state
func (p *ConnectedPlugin) GetState() PluginState {
p.mu.RLock()
defer p.mu.RUnlock()
return p.State
}
// SetState safely sets plugin state
func (p *ConnectedPlugin) SetState(state PluginState) {
p.mu.Lock()
defer p.mu.Unlock()
p.State = state
p.LastHeartbeat = time.Now()
}
// IsHealthy safely checks if plugin is healthy
func (p *ConnectedPlugin) IsHealthy() bool {
p.mu.RLock()
defer p.mu.RUnlock()
return p.Healthy
}
// UpdateHeartbeat safely updates heartbeat
func (p *ConnectedPlugin) UpdateHeartbeat(cpu, memory float32, pending, running int32) {
p.mu.Lock()
defer p.mu.Unlock()
p.LastHeartbeat = time.Now()
p.CPUUsagePercent = cpu
p.MemoryUsageMB = memory
p.PendingJobs = pending
p.RunningJobs = running
}
// GetConfig safely gets configuration
func (c *JobTypeConfig) GetConfig() (*plugin_pb.JobTypeConfig, error) {
c.mu.RLock()
defer c.mu.RUnlock()
return &plugin_pb.JobTypeConfig{
JobType: c.JobType,
Enabled: c.Enabled,
AdminConfig: c.AdminConfig,
WorkerConfig: c.WorkerConfig,
CreatedAt: timestamppb.New(c.CreatedAt),
UpdatedAt: timestamppb.New(c.UpdatedAt),
CreatedBy: c.CreatedBy,
}, nil
}
// SetConfig safely sets configuration
func (c *JobTypeConfig) SetConfig(enabled bool, adminConfig, workerConfig []*plugin_pb.ConfigFieldValue) {
c.mu.Lock()
defer c.mu.Unlock()
c.Enabled = enabled
c.AdminConfig = adminConfig
c.WorkerConfig = workerConfig
c.UpdatedAt = time.Now()
}
+1
View File
@@ -14,6 +14,7 @@ gen:
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 plugin.proto --go_out=./plugin_pb --go-grpc_out=./plugin_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
+558
View File
@@ -0,0 +1,558 @@
syntax = "proto3";
package seaweed_pb;
option go_package = "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb";
import "google/protobuf/duration.proto";
import "google/protobuf/timestamp.proto";
import "google/protobuf/empty.proto";
// PluginService: Worker-to-Admin communication service
// Workers establish bidirectional streaming connections to register and execute jobs
service PluginService {
// Connect: Bidirectional stream for plugin registration, heartbeat, and lifecycle
rpc Connect(stream PluginMessage) returns (stream AdminMessage);
// ExecuteJob: Bidirectional stream for job execution with progress updates
rpc ExecuteJob(stream JobExecutionMessage) returns (stream JobProgressMessage);
}
// AdminQueryService: Read-only queries for UI and monitoring
service AdminQueryService {
// GetPluginStats: Get current plugin system statistics
rpc GetPluginStats(google.protobuf.Empty) returns (PluginStats);
// ListPlugins: List all connected plugins with their capabilities
rpc ListPlugins(google.protobuf.Empty) returns (PluginList);
// ListJobs: List jobs by type and status
rpc ListJobs(ListJobsRequest) returns (JobList);
// GetJob: Get detailed job information
rpc GetJob(GetJobRequest) returns (JobDetail);
// GetJobExecutionLog: Stream job execution logs
rpc GetJobExecutionLog(GetJobExecutionLogRequest) returns (stream ExecutionLogEntry);
}
// AdminCommandService: Write operations for job management
service AdminCommandService {
// UpdateJobTypeConfig: Update configuration for a job type
rpc UpdateJobTypeConfig(UpdateJobTypeConfigRequest) returns (JobTypeConfig);
// CreateJob: Manually create a job
rpc CreateJob(CreateJobRequest) returns (Job);
// CancelJob: Cancel a running job
rpc CancelJob(CancelJobRequest) returns (google.protobuf.Empty);
// RetryJob: Retry a failed job
rpc RetryJob(RetryJobRequest) returns (Job);
}
// ============================================================================
// LIFECYCLE AND REGISTRATION MESSAGES
// ============================================================================
// PluginMessage: Messages sent by worker to admin
message PluginMessage {
oneof content {
PluginRegister register = 1;
PluginHeartbeat heartbeat = 2;
ExecutionStatusUpdate status_update = 3;
}
}
// AdminMessage: Messages sent by admin to worker
message AdminMessage {
oneof content {
JobRequest job_request = 1;
ConfigUpdate config_update = 2;
AdminCommand admin_command = 3;
}
}
// PluginRegister: Initial registration message from worker
message PluginRegister {
string plugin_id = 1; // Unique plugin identifier
string name = 2; // Human-readable plugin name
string version = 3; // Plugin version
string protocol_version = 4; // Protocol version (e.g., "v1")
repeated JobTypeCapability capabilities = 5; // Job types this plugin can handle
}
// JobTypeCapability: Declares what job type a plugin handles
message JobTypeCapability {
string job_type = 1; // Job type identifier (e.g., "erasure_coding", "vacuum")
bool can_detect = 2; // Plugin can detect jobs of this type
bool can_execute = 3; // Plugin can execute jobs of this type
string version = 4; // Job type version
}
// PluginHeartbeat: Periodic heartbeat to keep connection alive
message PluginHeartbeat {
string plugin_id = 1;
google.protobuf.Timestamp timestamp = 2;
int64 uptime_seconds = 3;
int32 pending_jobs = 4;
float cpu_usage_percent = 5;
float memory_usage_mb = 6;
}
// ============================================================================
// CONFIGURATION SCHEMA MESSAGES
// ============================================================================
// ConfigField: Declarative field definition for UI generation
message ConfigField {
enum FieldType {
BOOL = 0;
INT = 1;
FLOAT = 2;
STRING = 3;
DURATION = 4;
SELECT = 5;
MULTISELECT = 6;
SECRET = 7;
TEXTAREA = 8;
JSON = 9;
PERCENTAGE = 10;
BYTES = 11;
CRON = 12;
}
string name = 1; // Field identifier
string label = 2; // Display label
string description = 3; // Help text
FieldType field_type = 4; // UI control type
bool required = 5; // Is required
string default_value = 6; // Default value
repeated ValidationRule validation_rules = 7; // Validation constraints
repeated ConfigOption options = 8; // For SELECT/MULTISELECT
ConfigField.Options options_msg = 9; // Additional options
message Options {
string placeholder = 1;
repeated string suggestions = 2;
int32 min_length = 3;
int32 max_length = 4;
int32 min_value = 5;
int32 max_value = 6;
float min_float = 7;
float max_float = 8;
}
}
// ValidationRule: Constraint for field validation
message ValidationRule {
enum RuleType {
MIN_VALUE = 0;
MAX_VALUE = 1;
PATTERN = 2;
MIN_LENGTH = 3;
MAX_LENGTH = 4;
CUSTOM = 5;
}
RuleType rule_type = 1;
string value = 2;
string error_message = 3;
}
// ConfigOption: Selection option for SELECT/MULTISELECT
message ConfigOption {
string value = 1;
string label = 2;
string description = 3;
}
// ConfigFieldValue: Value for a config field
message ConfigFieldValue {
string field_name = 1;
string string_value = 2; // For STRING, CRON
int64 int_value = 3; // For INT, BYTES
float float_value = 4; // For FLOAT, PERCENTAGE
bool bool_value = 5; // For BOOL
google.protobuf.Duration duration_value = 6; // For DURATION
repeated string multiselect_values = 7; // For MULTISELECT
string json_value = 8; // For JSON
}
// JobTypeConfig: Configuration for a job type
message JobTypeConfig {
string job_type = 1;
bool enabled = 2;
repeated ConfigFieldValue admin_config = 3; // Admin-managed settings
repeated ConfigFieldValue worker_config = 4; // Worker-managed settings
google.protobuf.Timestamp created_at = 5;
google.protobuf.Timestamp updated_at = 6;
string created_by = 7;
}
// ConfigUpdateRequest: Request to update job type configuration
message ConfigUpdateRequest {
string job_type = 1;
repeated ConfigFieldValue config_values = 2;
}
// ============================================================================
// CONFIGURATION SCHEMA DISCOVERY
// ============================================================================
// GetConfigurationSchemaRequest: Request for config schema
message GetConfigurationSchemaRequest {
string job_type = 1;
}
// JobTypeConfigSchema: Complete configuration schema for a job type
message JobTypeConfigSchema {
string job_type = 1;
string version = 2;
string description = 3;
repeated ConfigField admin_fields = 4; // Fields managed by admin UI
repeated ConfigField worker_fields = 5; // Fields managed by worker
repeated ConfigGroup field_groups = 6; // Field grouping for UI
}
// ConfigGroup: Grouping of related fields for UI organization
message ConfigGroup {
string name = 1;
string label = 2;
string description = 3;
repeated string field_names = 4;
}
// ============================================================================
// JOB DETECTION MESSAGES
// ============================================================================
// DetectionRequest: Request to detect jobs
message DetectionRequest {
string job_type = 1;
JobTypeConfig config = 2; // Current configuration
repeated string filter_tags = 3; // Optional filtering
}
// DetectionResponse: Response from detection
message DetectionResponse {
repeated DetectedJob detected_jobs = 1;
google.protobuf.Timestamp detection_time = 2;
string detector_plugin_id = 3;
}
// DetectedJob: A job candidate detected by worker
message DetectedJob {
string job_key = 1; // Unique key for deduplication
string job_type = 2;
string description = 3; // Human-readable description
int64 priority = 4; // Higher = higher priority
google.protobuf.Duration estimated_duration = 5; // Estimated execution time
map<string, string> metadata = 6; // Job-specific metadata
repeated ConfigFieldValue suggested_config = 7; // Suggested overrides
}
// ============================================================================
// JOB EXECUTION MESSAGES
// ============================================================================
// JobRequest: Request to execute a job
message JobRequest {
string job_id = 1;
string job_type = 2;
string description = 3;
int64 priority = 4;
google.protobuf.Timestamp created_at = 5;
repeated ConfigFieldValue config = 6; // Job-specific config
map<string, string> metadata = 7; // Job metadata
string request_source = 8; // admin | detection | retry
}
// JobExecutionMessage: Messages from worker during job execution
message JobExecutionMessage {
string job_id = 1;
oneof content {
JobStarted job_started = 2;
JobProgress progress = 3;
JobCheckpoint checkpoint = 4;
JobCompleted job_completed = 5;
JobFailed job_failed = 6;
ExecutionLog log_entry = 7;
}
}
// JobStarted: Job execution started
message JobStarted {
google.protobuf.Timestamp started_at = 1;
string executor_id = 2; // Plugin executor ID
}
// JobProgress: Progress update during execution
message JobProgress {
int32 progress_percent = 1; // 0-100
string current_step = 2;
string status_message = 3;
google.protobuf.Timestamp updated_at = 4;
}
// JobCheckpoint: Checkpoint for resuming jobs
message JobCheckpoint {
string checkpoint_id = 1;
bytes checkpoint_data = 2; // Serialized checkpoint
google.protobuf.Timestamp created_at = 3;
}
// JobCompleted: Job execution completed successfully
message JobCompleted {
google.protobuf.Timestamp completed_at = 1;
string summary = 2;
map<string, string> output = 3; // Job output data
}
// JobFailed: Job execution failed
message JobFailed {
string error_code = 1; // Error classification
string error_message = 2; // Human-readable error
bytes error_details = 3; // Detailed error data
bool retryable = 4; // Can this job be retried
google.protobuf.Timestamp failed_at = 5;
int32 retry_count = 6;
}
// ExecutionLog: Log entry during job execution
message ExecutionLog {
enum LogLevel {
DEBUG = 0;
INFO = 1;
WARNING = 2;
ERROR = 3;
}
google.protobuf.Timestamp timestamp = 1;
LogLevel level = 2;
string message = 3;
map<string, string> context = 4;
}
// JobProgressMessage: Messages from admin to worker
message JobProgressMessage {
string job_id = 1;
oneof content {
JobProgressUpdate update = 2;
ExecutionCommand command = 3;
}
}
// JobProgressUpdate: Status update for UI
message JobProgressUpdate {
string job_id = 1;
int32 progress_percent = 2;
string current_step = 3;
string status_message = 4;
google.protobuf.Timestamp timestamp = 5;
}
// ExecutionCommand: Command from admin to worker
message ExecutionCommand {
enum CommandType {
PAUSE = 0;
RESUME = 1;
CANCEL = 2;
GET_STATUS = 3;
}
string job_id = 1;
CommandType command_type = 2;
map<string, string> parameters = 3;
}
// ============================================================================
// STATUS UPDATE MESSAGES
// ============================================================================
// ExecutionStatusUpdate: Status update from worker
message ExecutionStatusUpdate {
string plugin_id = 1;
string job_id = 2;
string status = 3; // running | completed | failed | paused
int32 progress_percent = 4;
google.protobuf.Timestamp timestamp = 5;
}
// ConfigUpdate: Configuration update message
message ConfigUpdate {
string job_type = 1;
repeated ConfigFieldValue config_values = 2;
google.protobuf.Timestamp updated_at = 3;
}
// AdminCommand: Admin command to worker
message AdminCommand {
enum CommandType {
RELOAD_CONFIG = 0;
ENABLE_JOB_TYPE = 1;
DISABLE_JOB_TYPE = 2;
SHUTDOWN = 3;
}
CommandType command_type = 1;
map<string, string> parameters = 2;
}
// ============================================================================
// QUERY/COMMAND REQUEST/RESPONSE MESSAGES
// ============================================================================
// ListJobsRequest: Request to list jobs
message ListJobsRequest {
string job_type = 1;
enum StatusFilter {
ALL = 0;
PENDING = 1;
RUNNING = 2;
COMPLETED = 3;
FAILED = 4;
}
StatusFilter status = 2;
int32 limit = 3;
int32 offset = 4;
}
// GetJobRequest: Request for specific job
message GetJobRequest {
string job_id = 1;
}
// GetJobExecutionLogRequest: Request for job logs
message GetJobExecutionLogRequest {
string job_id = 1;
enum LogLevel {
DEBUG = 0;
INFO = 1;
WARNING = 2;
ERROR = 3;
ALL = 4;
}
LogLevel level = 2;
int32 tail_lines = 3;
}
// ============================================================================
// RESPONSE MESSAGES
// ============================================================================
// Job: Job information
message Job {
string job_id = 1;
string job_type = 2;
string description = 3;
int64 priority = 4;
enum Status {
PENDING = 0;
RUNNING = 1;
COMPLETED = 2;
FAILED = 3;
CANCELLED = 4;
PAUSED = 5;
}
Status status = 5;
google.protobuf.Timestamp created_at = 6;
google.protobuf.Timestamp updated_at = 7;
string executor_plugin_id = 8;
int32 progress_percent = 9;
repeated ConfigFieldValue config = 10;
}
// JobDetail: Detailed job information with execution history
message JobDetail {
Job job = 1;
repeated ExecutionLogEntry log_entries = 2;
JobCompleted completion_info = 3;
JobFailed failure_info = 4;
repeated string checkpoint_ids = 5;
}
// ExecutionLogEntry: Log entry with metadata
message ExecutionLogEntry {
google.protobuf.Timestamp timestamp = 1;
ExecutionLog.LogLevel level = 2;
string message = 3;
map<string, string> context = 4;
}
// JobList: List of jobs
message JobList {
repeated Job jobs = 1;
int32 total_count = 2;
int32 returned_count = 3;
}
// PluginStats: System-wide statistics
message PluginStats {
int32 total_plugins = 1;
int32 active_plugins = 2;
int32 total_jobs = 3;
int32 pending_jobs = 4;
int32 running_jobs = 5;
int32 completed_jobs = 6;
int32 failed_jobs = 7;
google.protobuf.Timestamp last_update = 8;
map<string, JobTypeStats> job_type_stats = 9;
}
// JobTypeStats: Statistics for a job type
message JobTypeStats {
string job_type = 1;
int32 total = 2;
int32 pending = 3;
int32 running = 4;
int32 completed = 5;
int32 failed = 6;
float success_rate = 7;
google.protobuf.Duration avg_execution_time = 8;
}
// PluginList: List of connected plugins
message PluginList {
repeated PluginInfo plugins = 1;
}
// PluginInfo: Information about a connected plugin
message PluginInfo {
string plugin_id = 1;
string name = 2;
string version = 3;
string protocol_version = 4;
google.protobuf.Timestamp connected_at = 5;
google.protobuf.Timestamp last_heartbeat = 6;
repeated JobTypeCapability capabilities = 7;
bool healthy = 8;
string status = 9;
}
// UpdateJobTypeConfigRequest: Request to update config
message UpdateJobTypeConfigRequest {
JobTypeConfig config = 1;
}
// CreateJobRequest: Request to create a job
message CreateJobRequest {
string job_type = 1;
string description = 2;
int64 priority = 3;
repeated ConfigFieldValue config = 4;
map<string, string> metadata = 5;
}
// CancelJobRequest: Request to cancel a job
message CancelJobRequest {
string job_id = 1;
}
// RetryJobRequest: Request to retry a job
message RetryJobRequest {
string job_id = 1;
}
File diff suppressed because it is too large. Load diff
+658
View File
@@ -0,0 +1,658 @@
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
// versions:
// - protoc-gen-go-grpc v1.5.1
// - protoc v6.33.4
// source: plugin.proto
package plugin_pb
import (
context "context"
grpc "google.golang.org/grpc"
codes "google.golang.org/grpc/codes"
status "google.golang.org/grpc/status"
emptypb "google.golang.org/protobuf/types/known/emptypb"
)
// 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 (
PluginService_Connect_FullMethodName = "/seaweed_pb.PluginService/Connect"
PluginService_ExecuteJob_FullMethodName = "/seaweed_pb.PluginService/ExecuteJob"
)
// PluginServiceClient is the client API for PluginService 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.
//
// PluginService: Worker-to-Admin communication service
// Workers establish bidirectional streaming connections to register and execute jobs
type PluginServiceClient interface {
// Connect: Bidirectional stream for plugin registration, heartbeat, and lifecycle
Connect(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[PluginMessage, AdminMessage], error)
// ExecuteJob: Bidirectional stream for job execution with progress updates
ExecuteJob(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[JobExecutionMessage, JobProgressMessage], error)
}
type pluginServiceClient struct {
cc grpc.ClientConnInterface
}
func NewPluginServiceClient(cc grpc.ClientConnInterface) PluginServiceClient {
return &pluginServiceClient{cc}
}
func (c *pluginServiceClient) Connect(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[PluginMessage, AdminMessage], error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
stream, err := c.cc.NewStream(ctx, &PluginService_ServiceDesc.Streams[0], PluginService_Connect_FullMethodName, cOpts...)
if err != nil {
return nil, err
}
x := &grpc.GenericClientStream[PluginMessage, 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 PluginService_ConnectClient = grpc.BidiStreamingClient[PluginMessage, AdminMessage]
func (c *pluginServiceClient) ExecuteJob(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[JobExecutionMessage, JobProgressMessage], error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
stream, err := c.cc.NewStream(ctx, &PluginService_ServiceDesc.Streams[1], PluginService_ExecuteJob_FullMethodName, cOpts...)
if err != nil {
return nil, err
}
x := &grpc.GenericClientStream[JobExecutionMessage, JobProgressMessage]{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 PluginService_ExecuteJobClient = grpc.BidiStreamingClient[JobExecutionMessage, JobProgressMessage]
// PluginServiceServer is the server API for PluginService service.
// All implementations must embed UnimplementedPluginServiceServer
// for forward compatibility.
//
// PluginService: Worker-to-Admin communication service
// Workers establish bidirectional streaming connections to register and execute jobs
type PluginServiceServer interface {
// Connect: Bidirectional stream for plugin registration, heartbeat, and lifecycle
Connect(grpc.BidiStreamingServer[PluginMessage, AdminMessage]) error
// ExecuteJob: Bidirectional stream for job execution with progress updates
ExecuteJob(grpc.BidiStreamingServer[JobExecutionMessage, JobProgressMessage]) error
mustEmbedUnimplementedPluginServiceServer()
}
// UnimplementedPluginServiceServer 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 UnimplementedPluginServiceServer struct{}
func (UnimplementedPluginServiceServer) Connect(grpc.BidiStreamingServer[PluginMessage, AdminMessage]) error {
return status.Errorf(codes.Unimplemented, "method Connect not implemented")
}
func (UnimplementedPluginServiceServer) ExecuteJob(grpc.BidiStreamingServer[JobExecutionMessage, JobProgressMessage]) error {
return status.Errorf(codes.Unimplemented, "method ExecuteJob not implemented")
}
func (UnimplementedPluginServiceServer) mustEmbedUnimplementedPluginServiceServer() {}
func (UnimplementedPluginServiceServer) testEmbeddedByValue() {}
// UnsafePluginServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to PluginServiceServer will
// result in compilation errors.
type UnsafePluginServiceServer interface {
mustEmbedUnimplementedPluginServiceServer()
}
func RegisterPluginServiceServer(s grpc.ServiceRegistrar, srv PluginServiceServer) {
// If the following call pancis, it indicates UnimplementedPluginServiceServer 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(&PluginService_ServiceDesc, srv)
}
func _PluginService_Connect_Handler(srv interface{}, stream grpc.ServerStream) error {
return srv.(PluginServiceServer).Connect(&grpc.GenericServerStream[PluginMessage, 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 PluginService_ConnectServer = grpc.BidiStreamingServer[PluginMessage, AdminMessage]
func _PluginService_ExecuteJob_Handler(srv interface{}, stream grpc.ServerStream) error {
return srv.(PluginServiceServer).ExecuteJob(&grpc.GenericServerStream[JobExecutionMessage, JobProgressMessage]{ServerStream: stream})
}
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
type PluginService_ExecuteJobServer = grpc.BidiStreamingServer[JobExecutionMessage, JobProgressMessage]
// PluginService_ServiceDesc is the grpc.ServiceDesc for PluginService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var PluginService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "seaweed_pb.PluginService",
HandlerType: (*PluginServiceServer)(nil),
Methods: []grpc.MethodDesc{},
Streams: []grpc.StreamDesc{
{
StreamName: "Connect",
Handler: _PluginService_Connect_Handler,
ServerStreams: true,
ClientStreams: true,
},
{
StreamName: "ExecuteJob",
Handler: _PluginService_ExecuteJob_Handler,
ServerStreams: true,
ClientStreams: true,
},
},
Metadata: "plugin.proto",
}
const (
AdminQueryService_GetPluginStats_FullMethodName = "/seaweed_pb.AdminQueryService/GetPluginStats"
AdminQueryService_ListPlugins_FullMethodName = "/seaweed_pb.AdminQueryService/ListPlugins"
AdminQueryService_ListJobs_FullMethodName = "/seaweed_pb.AdminQueryService/ListJobs"
AdminQueryService_GetJob_FullMethodName = "/seaweed_pb.AdminQueryService/GetJob"
AdminQueryService_GetJobExecutionLog_FullMethodName = "/seaweed_pb.AdminQueryService/GetJobExecutionLog"
)
// AdminQueryServiceClient is the client API for AdminQueryService 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.
//
// AdminQueryService: Read-only queries for UI and monitoring
type AdminQueryServiceClient interface {
// GetPluginStats: Get current plugin system statistics
GetPluginStats(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*PluginStats, error)
// ListPlugins: List all connected plugins with their capabilities
ListPlugins(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*PluginList, error)
// ListJobs: List jobs by type and status
ListJobs(ctx context.Context, in *ListJobsRequest, opts ...grpc.CallOption) (*JobList, error)
// GetJob: Get detailed job information
GetJob(ctx context.Context, in *GetJobRequest, opts ...grpc.CallOption) (*JobDetail, error)
// GetJobExecutionLog: Stream job execution logs
GetJobExecutionLog(ctx context.Context, in *GetJobExecutionLogRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[ExecutionLogEntry], error)
}
type adminQueryServiceClient struct {
cc grpc.ClientConnInterface
}
func NewAdminQueryServiceClient(cc grpc.ClientConnInterface) AdminQueryServiceClient {
return &adminQueryServiceClient{cc}
}
func (c *adminQueryServiceClient) GetPluginStats(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*PluginStats, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(PluginStats)
err := c.cc.Invoke(ctx, AdminQueryService_GetPluginStats_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminQueryServiceClient) ListPlugins(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*PluginList, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(PluginList)
err := c.cc.Invoke(ctx, AdminQueryService_ListPlugins_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminQueryServiceClient) ListJobs(ctx context.Context, in *ListJobsRequest, opts ...grpc.CallOption) (*JobList, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(JobList)
err := c.cc.Invoke(ctx, AdminQueryService_ListJobs_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminQueryServiceClient) GetJob(ctx context.Context, in *GetJobRequest, opts ...grpc.CallOption) (*JobDetail, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(JobDetail)
err := c.cc.Invoke(ctx, AdminQueryService_GetJob_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminQueryServiceClient) GetJobExecutionLog(ctx context.Context, in *GetJobExecutionLogRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[ExecutionLogEntry], error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
stream, err := c.cc.NewStream(ctx, &AdminQueryService_ServiceDesc.Streams[0], AdminQueryService_GetJobExecutionLog_FullMethodName, cOpts...)
if err != nil {
return nil, err
}
x := &grpc.GenericClientStream[GetJobExecutionLogRequest, ExecutionLogEntry]{ClientStream: stream}
if err := x.ClientStream.SendMsg(in); err != nil {
return nil, err
}
if err := x.ClientStream.CloseSend(); err != nil {
return nil, err
}
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 AdminQueryService_GetJobExecutionLogClient = grpc.ServerStreamingClient[ExecutionLogEntry]
// AdminQueryServiceServer is the server API for AdminQueryService service.
// All implementations must embed UnimplementedAdminQueryServiceServer
// for forward compatibility.
//
// AdminQueryService: Read-only queries for UI and monitoring
type AdminQueryServiceServer interface {
// GetPluginStats: Get current plugin system statistics
GetPluginStats(context.Context, *emptypb.Empty) (*PluginStats, error)
// ListPlugins: List all connected plugins with their capabilities
ListPlugins(context.Context, *emptypb.Empty) (*PluginList, error)
// ListJobs: List jobs by type and status
ListJobs(context.Context, *ListJobsRequest) (*JobList, error)
// GetJob: Get detailed job information
GetJob(context.Context, *GetJobRequest) (*JobDetail, error)
// GetJobExecutionLog: Stream job execution logs
GetJobExecutionLog(*GetJobExecutionLogRequest, grpc.ServerStreamingServer[ExecutionLogEntry]) error
mustEmbedUnimplementedAdminQueryServiceServer()
}
// UnimplementedAdminQueryServiceServer 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 UnimplementedAdminQueryServiceServer struct{}
func (UnimplementedAdminQueryServiceServer) GetPluginStats(context.Context, *emptypb.Empty) (*PluginStats, error) {
return nil, status.Errorf(codes.Unimplemented, "method GetPluginStats not implemented")
}
func (UnimplementedAdminQueryServiceServer) ListPlugins(context.Context, *emptypb.Empty) (*PluginList, error) {
return nil, status.Errorf(codes.Unimplemented, "method ListPlugins not implemented")
}
func (UnimplementedAdminQueryServiceServer) ListJobs(context.Context, *ListJobsRequest) (*JobList, error) {
return nil, status.Errorf(codes.Unimplemented, "method ListJobs not implemented")
}
func (UnimplementedAdminQueryServiceServer) GetJob(context.Context, *GetJobRequest) (*JobDetail, error) {
return nil, status.Errorf(codes.Unimplemented, "method GetJob not implemented")
}
func (UnimplementedAdminQueryServiceServer) GetJobExecutionLog(*GetJobExecutionLogRequest, grpc.ServerStreamingServer[ExecutionLogEntry]) error {
return status.Errorf(codes.Unimplemented, "method GetJobExecutionLog not implemented")
}
func (UnimplementedAdminQueryServiceServer) mustEmbedUnimplementedAdminQueryServiceServer() {}
func (UnimplementedAdminQueryServiceServer) testEmbeddedByValue() {}
// UnsafeAdminQueryServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to AdminQueryServiceServer will
// result in compilation errors.
type UnsafeAdminQueryServiceServer interface {
mustEmbedUnimplementedAdminQueryServiceServer()
}
func RegisterAdminQueryServiceServer(s grpc.ServiceRegistrar, srv AdminQueryServiceServer) {
// If the following call pancis, it indicates UnimplementedAdminQueryServiceServer 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(&AdminQueryService_ServiceDesc, srv)
}
func _AdminQueryService_GetPluginStats_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(emptypb.Empty)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminQueryServiceServer).GetPluginStats(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminQueryService_GetPluginStats_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminQueryServiceServer).GetPluginStats(ctx, req.(*emptypb.Empty))
}
return interceptor(ctx, in, info, handler)
}
func _AdminQueryService_ListPlugins_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(emptypb.Empty)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminQueryServiceServer).ListPlugins(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminQueryService_ListPlugins_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminQueryServiceServer).ListPlugins(ctx, req.(*emptypb.Empty))
}
return interceptor(ctx, in, info, handler)
}
func _AdminQueryService_ListJobs_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(ListJobsRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminQueryServiceServer).ListJobs(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminQueryService_ListJobs_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminQueryServiceServer).ListJobs(ctx, req.(*ListJobsRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AdminQueryService_GetJob_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(GetJobRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminQueryServiceServer).GetJob(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminQueryService_GetJob_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminQueryServiceServer).GetJob(ctx, req.(*GetJobRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AdminQueryService_GetJobExecutionLog_Handler(srv interface{}, stream grpc.ServerStream) error {
m := new(GetJobExecutionLogRequest)
if err := stream.RecvMsg(m); err != nil {
return err
}
return srv.(AdminQueryServiceServer).GetJobExecutionLog(m, &grpc.GenericServerStream[GetJobExecutionLogRequest, ExecutionLogEntry]{ServerStream: stream})
}
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
type AdminQueryService_GetJobExecutionLogServer = grpc.ServerStreamingServer[ExecutionLogEntry]
// AdminQueryService_ServiceDesc is the grpc.ServiceDesc for AdminQueryService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var AdminQueryService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "seaweed_pb.AdminQueryService",
HandlerType: (*AdminQueryServiceServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "GetPluginStats",
Handler: _AdminQueryService_GetPluginStats_Handler,
},
{
MethodName: "ListPlugins",
Handler: _AdminQueryService_ListPlugins_Handler,
},
{
MethodName: "ListJobs",
Handler: _AdminQueryService_ListJobs_Handler,
},
{
MethodName: "GetJob",
Handler: _AdminQueryService_GetJob_Handler,
},
},
Streams: []grpc.StreamDesc{
{
StreamName: "GetJobExecutionLog",
Handler: _AdminQueryService_GetJobExecutionLog_Handler,
ServerStreams: true,
},
},
Metadata: "plugin.proto",
}
const (
AdminCommandService_UpdateJobTypeConfig_FullMethodName = "/seaweed_pb.AdminCommandService/UpdateJobTypeConfig"
AdminCommandService_CreateJob_FullMethodName = "/seaweed_pb.AdminCommandService/CreateJob"
AdminCommandService_CancelJob_FullMethodName = "/seaweed_pb.AdminCommandService/CancelJob"
AdminCommandService_RetryJob_FullMethodName = "/seaweed_pb.AdminCommandService/RetryJob"
)
// AdminCommandServiceClient is the client API for AdminCommandService 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.
//
// AdminCommandService: Write operations for job management
type AdminCommandServiceClient interface {
// UpdateJobTypeConfig: Update configuration for a job type
UpdateJobTypeConfig(ctx context.Context, in *UpdateJobTypeConfigRequest, opts ...grpc.CallOption) (*JobTypeConfig, error)
// CreateJob: Manually create a job
CreateJob(ctx context.Context, in *CreateJobRequest, opts ...grpc.CallOption) (*Job, error)
// CancelJob: Cancel a running job
CancelJob(ctx context.Context, in *CancelJobRequest, opts ...grpc.CallOption) (*emptypb.Empty, error)
// RetryJob: Retry a failed job
RetryJob(ctx context.Context, in *RetryJobRequest, opts ...grpc.CallOption) (*Job, error)
}
type adminCommandServiceClient struct {
cc grpc.ClientConnInterface
}
func NewAdminCommandServiceClient(cc grpc.ClientConnInterface) AdminCommandServiceClient {
return &adminCommandServiceClient{cc}
}
func (c *adminCommandServiceClient) UpdateJobTypeConfig(ctx context.Context, in *UpdateJobTypeConfigRequest, opts ...grpc.CallOption) (*JobTypeConfig, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(JobTypeConfig)
err := c.cc.Invoke(ctx, AdminCommandService_UpdateJobTypeConfig_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminCommandServiceClient) CreateJob(ctx context.Context, in *CreateJobRequest, opts ...grpc.CallOption) (*Job, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(Job)
err := c.cc.Invoke(ctx, AdminCommandService_CreateJob_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminCommandServiceClient) CancelJob(ctx context.Context, in *CancelJobRequest, opts ...grpc.CallOption) (*emptypb.Empty, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(emptypb.Empty)
err := c.cc.Invoke(ctx, AdminCommandService_CancelJob_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *adminCommandServiceClient) RetryJob(ctx context.Context, in *RetryJobRequest, opts ...grpc.CallOption) (*Job, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(Job)
err := c.cc.Invoke(ctx, AdminCommandService_RetryJob_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
// AdminCommandServiceServer is the server API for AdminCommandService service.
// All implementations must embed UnimplementedAdminCommandServiceServer
// for forward compatibility.
//
// AdminCommandService: Write operations for job management
type AdminCommandServiceServer interface {
// UpdateJobTypeConfig: Update configuration for a job type
UpdateJobTypeConfig(context.Context, *UpdateJobTypeConfigRequest) (*JobTypeConfig, error)
// CreateJob: Manually create a job
CreateJob(context.Context, *CreateJobRequest) (*Job, error)
// CancelJob: Cancel a running job
CancelJob(context.Context, *CancelJobRequest) (*emptypb.Empty, error)
// RetryJob: Retry a failed job
RetryJob(context.Context, *RetryJobRequest) (*Job, error)
mustEmbedUnimplementedAdminCommandServiceServer()
}
// UnimplementedAdminCommandServiceServer 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 UnimplementedAdminCommandServiceServer struct{}
func (UnimplementedAdminCommandServiceServer) UpdateJobTypeConfig(context.Context, *UpdateJobTypeConfigRequest) (*JobTypeConfig, error) {
return nil, status.Errorf(codes.Unimplemented, "method UpdateJobTypeConfig not implemented")
}
func (UnimplementedAdminCommandServiceServer) CreateJob(context.Context, *CreateJobRequest) (*Job, error) {
return nil, status.Errorf(codes.Unimplemented, "method CreateJob not implemented")
}
func (UnimplementedAdminCommandServiceServer) CancelJob(context.Context, *CancelJobRequest) (*emptypb.Empty, error) {
return nil, status.Errorf(codes.Unimplemented, "method CancelJob not implemented")
}
func (UnimplementedAdminCommandServiceServer) RetryJob(context.Context, *RetryJobRequest) (*Job, error) {
return nil, status.Errorf(codes.Unimplemented, "method RetryJob not implemented")
}
func (UnimplementedAdminCommandServiceServer) mustEmbedUnimplementedAdminCommandServiceServer() {}
func (UnimplementedAdminCommandServiceServer) testEmbeddedByValue() {}
// UnsafeAdminCommandServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to AdminCommandServiceServer will
// result in compilation errors.
type UnsafeAdminCommandServiceServer interface {
mustEmbedUnimplementedAdminCommandServiceServer()
}
func RegisterAdminCommandServiceServer(s grpc.ServiceRegistrar, srv AdminCommandServiceServer) {
// If the following call pancis, it indicates UnimplementedAdminCommandServiceServer 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(&AdminCommandService_ServiceDesc, srv)
}
func _AdminCommandService_UpdateJobTypeConfig_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(UpdateJobTypeConfigRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminCommandServiceServer).UpdateJobTypeConfig(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminCommandService_UpdateJobTypeConfig_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminCommandServiceServer).UpdateJobTypeConfig(ctx, req.(*UpdateJobTypeConfigRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AdminCommandService_CreateJob_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(CreateJobRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminCommandServiceServer).CreateJob(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminCommandService_CreateJob_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminCommandServiceServer).CreateJob(ctx, req.(*CreateJobRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AdminCommandService_CancelJob_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(CancelJobRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminCommandServiceServer).CancelJob(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminCommandService_CancelJob_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminCommandServiceServer).CancelJob(ctx, req.(*CancelJobRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AdminCommandService_RetryJob_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(RetryJobRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AdminCommandServiceServer).RetryJob(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AdminCommandService_RetryJob_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AdminCommandServiceServer).RetryJob(ctx, req.(*RetryJobRequest))
}
return interceptor(ctx, in, info, handler)
}
// AdminCommandService_ServiceDesc is the grpc.ServiceDesc for AdminCommandService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var AdminCommandService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "seaweed_pb.AdminCommandService",
HandlerType: (*AdminCommandServiceServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "UpdateJobTypeConfig",
Handler: _AdminCommandService_UpdateJobTypeConfig_Handler,
},
{
MethodName: "CreateJob",
Handler: _AdminCommandService_CreateJob_Handler,
},
{
MethodName: "CancelJob",
Handler: _AdminCommandService_CancelJob_Handler,
},
{
MethodName: "RetryJob",
Handler: _AdminCommandService_RetryJob_Handler,
},
},
Streams: []grpc.StreamDesc{},
Metadata: "plugin.proto",
}