324 lines
11 KiB
Go
324 lines
11 KiB
Go
package dispatcher
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/configread"
|
|
"git.ipao.vip/rogee/go-sip/internal/mq"
|
|
"git.ipao.vip/rogee/go-sip/internal/store"
|
|
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
|
)
|
|
|
|
// CurrentRuntime owns task/control consumers for one Dispatcher. SaaS creates
|
|
// every queue/binding; this process only checks and consumes predeclared ones.
|
|
type CurrentRuntime struct {
|
|
Broker *mq.CurrentBroker
|
|
Bootstrap CurrentBootstrap
|
|
Execute CurrentExecuteController
|
|
Control CurrentControlController
|
|
PollInterval time.Duration
|
|
DiscoveryInterval time.Duration
|
|
Logger *slog.Logger
|
|
|
|
gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction
|
|
locksMu sync.Mutex
|
|
taskLocks map[string]*sync.Mutex
|
|
failures chan error
|
|
}
|
|
|
|
// Serve closes admission on every shutdown/failure; it never clears durable
|
|
// calls, results, task queues, or the SQLite file.
|
|
func (r *CurrentRuntime) Serve(ctx context.Context) (result error) {
|
|
if r == nil || r.Broker == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Bootstrap.DispatcherID == "" || r.Execute.Store != r.Bootstrap.Store || r.Control.Store != r.Bootstrap.Store || r.Execute.DispatcherID != r.Bootstrap.DispatcherID || r.Control.DispatcherID != r.Bootstrap.DispatcherID || r.PollInterval <= 0 || r.DiscoveryInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil {
|
|
return errors.New("current runtime requires one Dispatcher, durable state, verified SIP, independent clocks, and configured polling; external control drain is forbidden")
|
|
}
|
|
if err := r.Execute.validate(); err != nil {
|
|
return err
|
|
}
|
|
if r.Control.Client != r.Bootstrap.Client {
|
|
return errors.New("current runtime control and bootstrap must share the approved HTTP client")
|
|
}
|
|
r.failures = make(chan error, 1)
|
|
r.taskLocks = make(map[string]*sync.Mutex)
|
|
consumers := make(map[string]*mq.CurrentConsumer)
|
|
var controlConsumer *mq.CurrentConsumer
|
|
defer func() {
|
|
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
for queue, consumer := range consumers {
|
|
if err := consumer.Stop(stopCtx); err != nil {
|
|
result = errors.Join(result, fmt.Errorf("stop task consumer %q: %w", queue, err))
|
|
}
|
|
}
|
|
if controlConsumer != nil {
|
|
if err := controlConsumer.Stop(stopCtx); err != nil {
|
|
result = errors.Join(result, fmt.Errorf("stop control consumer: %w", err))
|
|
}
|
|
}
|
|
if err := r.Bootstrap.Store.CloseAdmission(r.Bootstrap.DispatcherID); err != nil {
|
|
result = errors.Join(result, fmt.Errorf("close Dispatcher admission: %w", err))
|
|
}
|
|
}()
|
|
|
|
var cursor string
|
|
var sip configread.CurrentSIP
|
|
r.Bootstrap.Cursor = &cursor
|
|
r.Bootstrap.SIP = &sip
|
|
r.Bootstrap.DrainControls = func(ctx context.Context) error {
|
|
_, err := r.Broker.DrainControlPredeclared(ctx, r.Broker.ControlQueue(), r.handleControl)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Start control consumption before bootstrap opens task admission.
|
|
consumer, err := r.Broker.StartPredeclaredConsumer(ctx, r.Broker.ControlQueue(), r.handleControl)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
controlConsumer = consumer
|
|
go r.watchConsumer(ctx, "control", consumer)
|
|
if err := r.Execute.FlushOutbox(ctx); err != nil {
|
|
return fmt.Errorf("recover durable results before admission: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
if err := r.Bootstrap.Run(ctx); err != nil {
|
|
if !errors.Is(err, store.ErrCurrentSIPPending) {
|
|
return fmt.Errorf("bootstrap current Dispatcher: %w", err)
|
|
}
|
|
r.Logger.Warn("control and outbox stay active while newer SIP revision waits; task admission remains closed", "dispatcher_id", r.Bootstrap.DispatcherID, "error", err)
|
|
}
|
|
follower := &CurrentDiscoveryFollower{DispatcherID: r.Bootstrap.DispatcherID, Client: r.Bootstrap.Client, Store: r.Bootstrap.Store, ApprovedSIP: sip, VerifySIP: r.Bootstrap.VerifySIP, Cursor: cursor}
|
|
if err := r.syncTaskConsumers(ctx, consumers); err != nil {
|
|
return err
|
|
}
|
|
if err := r.Execute.FlushOutbox(ctx); err != nil {
|
|
return err
|
|
}
|
|
poll := time.NewTicker(r.PollInterval)
|
|
defer poll.Stop()
|
|
discovery := time.NewTicker(r.DiscoveryInterval)
|
|
defer discovery.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case err := <-r.failures:
|
|
return fmt.Errorf("current MQ processing failed: %w", err)
|
|
case <-poll.C:
|
|
if err := r.processPending(ctx); err != nil {
|
|
return err
|
|
}
|
|
if err := r.Execute.FlushOutbox(ctx); err != nil {
|
|
return fmt.Errorf("deliver persisted results: %w", err)
|
|
}
|
|
if err := r.syncTaskConsumers(ctx, consumers); err != nil {
|
|
return err
|
|
}
|
|
case <-discovery.C:
|
|
if err := r.refreshSIP(ctx, follower); err != nil {
|
|
return fmt.Errorf("refresh approved SIP: %w", err)
|
|
}
|
|
_, pending, err := r.Bootstrap.Store.SIPState(r.Bootstrap.DispatcherID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if pending == 0 {
|
|
if err := follower.Poll(ctx); err != nil {
|
|
return fmt.Errorf("poll assigned tasks: %w", err)
|
|
}
|
|
}
|
|
if err := r.syncTaskConsumers(ctx, consumers); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (r *CurrentRuntime) watchConsumer(ctx context.Context, queue string, consumer *mq.CurrentConsumer) {
|
|
err := consumer.Wait(ctx)
|
|
if err != nil && ctx.Err() == nil {
|
|
r.Logger.Error("current MQ consumer failed", "dispatcher_id", r.Bootstrap.DispatcherID, "queue", queue, "error", err)
|
|
select {
|
|
case r.failures <- fmt.Errorf("queue %q: %w", queue, err):
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
|
|
// signalFailure forces the runtime to close admission after a consumer
|
|
// handler fails; RabbitMQ requeue alone would otherwise spin indefinitely.
|
|
func (r *CurrentRuntime) signalFailure(err error) {
|
|
if r.failures != nil {
|
|
select {
|
|
case r.failures <- err:
|
|
default: // the first failure is already being handled
|
|
}
|
|
}
|
|
}
|
|
|
|
func (r *CurrentRuntime) withTask(taskID string, process func() error) error {
|
|
r.gate.RLock()
|
|
defer r.gate.RUnlock()
|
|
r.locksMu.Lock()
|
|
lock := r.taskLocks[taskID]
|
|
if lock == nil {
|
|
lock = new(sync.Mutex)
|
|
r.taskLocks[taskID] = lock
|
|
}
|
|
r.locksMu.Unlock()
|
|
lock.Lock()
|
|
defer lock.Unlock()
|
|
return process()
|
|
}
|
|
|
|
func (r *CurrentRuntime) handleControl(ctx context.Context, _ string, body []byte) error {
|
|
var event struct {
|
|
EventType string `json:"event_type"`
|
|
Payload struct {
|
|
TaskID string `json:"task_id"`
|
|
Revision int64 `json:"revision"`
|
|
} `json:"payload"`
|
|
}
|
|
if err := json.Unmarshal(body, &event); err != nil {
|
|
failure := fmt.Errorf("decode MQ control routing: %w", err)
|
|
r.signalFailure(failure)
|
|
return failure
|
|
}
|
|
switch event.EventType {
|
|
case "task.control":
|
|
return r.withTask(event.Payload.TaskID, func() error {
|
|
if err := r.Control.ProcessControl(ctx, body); err != nil {
|
|
r.Logger.Error("task control failed", "dispatcher_id", r.Bootstrap.DispatcherID, "task_id", event.Payload.TaskID, "error", err)
|
|
r.signalFailure(err)
|
|
return err
|
|
}
|
|
return nil
|
|
})
|
|
case "sip.config":
|
|
r.gate.Lock()
|
|
err := r.Bootstrap.Store.NoteSIPChange(r.Bootstrap.DispatcherID, event.Payload.Revision)
|
|
r.gate.Unlock()
|
|
if err != nil {
|
|
failure := fmt.Errorf("persist SIP change before MQ ACK: %w", err)
|
|
r.signalFailure(failure)
|
|
return failure
|
|
}
|
|
r.Logger.Info("SIP notification persisted; task admission fenced until drain and loaded revision check", "dispatcher_id", r.Bootstrap.DispatcherID, "revision", event.Payload.Revision)
|
|
return nil
|
|
default:
|
|
failure := fmt.Errorf("unexpected MQ control event %q", event.EventType)
|
|
r.signalFailure(failure)
|
|
return failure
|
|
}
|
|
}
|
|
|
|
func (r *CurrentRuntime) handleTask(ctx context.Context, taskID string, body []byte) error {
|
|
return r.withTask(taskID, func() error {
|
|
if err := r.Execute.ProcessExecute(ctx, body); err != nil {
|
|
r.Logger.Error("call instruction failed", "dispatcher_id", r.Bootstrap.DispatcherID, "task_id", taskID, "error", err)
|
|
r.signalFailure(err)
|
|
return err
|
|
}
|
|
return nil
|
|
})
|
|
}
|
|
|
|
func (r *CurrentRuntime) processPending(ctx context.Context) error {
|
|
pending, err := r.Bootstrap.Store.ListPendingExecute(r.Bootstrap.DispatcherID)
|
|
if err != nil {
|
|
return fmt.Errorf("read durable rule-wait instructions: %w", err)
|
|
}
|
|
for _, command := range pending {
|
|
cmd := command
|
|
if err := r.withTask(cmd.TaskID, func() error { return r.Execute.dispatchPending(ctx, cmd) }); err != nil {
|
|
return fmt.Errorf("retry eligible pending instruction %q: %w", cmd.EventID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.CurrentConsumer) error {
|
|
assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
wanted := make(map[string]store.CurrentAssignedTask, len(assigned))
|
|
for _, task := range assigned {
|
|
route, err := tenant.CurrentTaskRoute(r.Bootstrap.DispatcherID, task.TaskID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" {
|
|
continue
|
|
}
|
|
if task.ControlState != "stopped" && task.ControlState != "stopping" && task.Status != "stopped" {
|
|
admitted, err := r.Bootstrap.Store.CanAdmit(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !admitted {
|
|
continue
|
|
}
|
|
waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if waiting != 0 {
|
|
continue
|
|
}
|
|
}
|
|
wanted[route.Queue] = task
|
|
}
|
|
for queue, consumer := range consumers {
|
|
if _, ok := wanted[queue]; ok {
|
|
continue
|
|
}
|
|
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
err := consumer.Stop(stopCtx)
|
|
cancel()
|
|
if err != nil {
|
|
return fmt.Errorf("stop nonadmitted task queue %q: %w", queue, err)
|
|
}
|
|
delete(consumers, queue)
|
|
}
|
|
for queue, task := range wanted {
|
|
if _, ok := consumers[queue]; ok {
|
|
continue
|
|
}
|
|
id, tenantID := task.TaskID, task.TenantID
|
|
ready := make(chan struct{})
|
|
var consumer *mq.CurrentConsumer
|
|
consumer, err = r.Broker.StartPredeclaredConsumer(ctx, queue, func(ctx context.Context, _ string, body []byte) error {
|
|
<-ready // the subscription is assigned before its first delivery can pause itself
|
|
if err := r.handleTask(ctx, id, body); err != nil {
|
|
return err
|
|
}
|
|
waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, tenantID, id)
|
|
if err != nil {
|
|
return fmt.Errorf("check task-local rule wait: %w", err)
|
|
}
|
|
if waiting != 0 {
|
|
consumer.RequestStop()
|
|
}
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("start SaaS-owned task queue %q: %w", queue, err)
|
|
}
|
|
consumers[queue] = consumer
|
|
close(ready)
|
|
go r.watchConsumer(ctx, queue, consumer)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Compile-time interface checks: the same broker confirms bound persistent
|
|
// results and consumes SaaS-owned queues without configure permissions.
|
|
var _ CurrentPublisher = (*mq.CurrentBroker)(nil)
|