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" ) // Runtime owns task/control consumers for one Dispatcher. SaaS creates // every queue/binding; this process checks and consumes predeclared ones and // purges a stopped task's Ready backlog without changing queue topology. type Runtime struct { Broker *mq.Broker Bootstrap Bootstrap Execute ExecuteController Control ControlController PollInterval time.Duration Logger *slog.Logger gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction consumersMu sync.Mutex // stop consumer and purge cannot race consumer sync taskConsumers map[string]*mq.Consumer locksMu sync.Mutex taskLocks map[string]*sync.Mutex failures chan error sipUpdates chan struct{} } // Serve closes admission on every shutdown/failure; it never clears durable // calls, results, or the SQLite file. Stop purges only the task's Ready backlog. func (r *Runtime) 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.Logger == nil || r.Bootstrap.DrainControls != nil { return errors.New("runtime requires one Dispatcher, durable state, verified SIP and configured pending-work 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.sipUpdates = make(chan struct{}, 1) r.taskLocks = make(map[string]*sync.Mutex) consumers := make(map[string]*mq.Consumer) r.taskConsumers = consumers r.Control.PurgeTaskQueue = r.Broker.PurgeTaskQueue var controlConsumer *mq.Consumer defer func() { stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if controlConsumer != nil { if err := controlConsumer.Stop(stopCtx); err != nil { result = errors.Join(result, fmt.Errorf("stop control consumer: %w", err)) } } r.consumersMu.Lock() 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)) } } r.consumersMu.Unlock() if err := r.Bootstrap.Store.CloseAdmission(r.Bootstrap.DispatcherID); err != nil { result = errors.Join(result, fmt.Errorf("close Dispatcher admission: %w", err)) } }() var sip configread.SIP var providers map[string]configread.Provider r.Bootstrap.SIP = &sip r.Bootstrap.Providers = &providers r.Control.ApprovedSIP = &sip r.Control.ApprovedProviders = &providers 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.ErrSIPPending) { 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) } 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() 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.refreshSIP(ctx, &sip); err != nil { return fmt.Errorf("apply pending SIP change: %w", err) } 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 <-r.sipUpdates: if err := r.refreshSIP(ctx, &sip); err != nil { return fmt.Errorf("apply SIP notification: %w", err) } } } } func (r *Runtime) watchConsumer(ctx context.Context, queue string, consumer *mq.Consumer) { 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 *Runtime) signalFailure(err error) { if r.failures != nil { select { case r.failures <- err: default: // the first failure is already being handled } } } func (r *Runtime) 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 *Runtime) handleControl(ctx context.Context, _ string, body []byte) error { var event struct { EventType string `json:"event_type"` Payload struct { TaskID string `json:"task_id"` Action string `json:"action"` 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": if event.Payload.Action == "stop" { // Stop the consumer before purging: closing its channel requeues any // unacknowledged deliveries so the purge covers them as Ready. // Hold this lock through the control transition to prevent sync from // restarting the consumer before the stop barrier is persisted. r.consumersMu.Lock() defer r.consumersMu.Unlock() route, err := tenant.TaskRoute(r.Bootstrap.DispatcherID, event.Payload.TaskID) if err != nil { r.signalFailure(err) return err } if consumer := r.taskConsumers[route.Queue]; consumer != nil { stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) err := consumer.Stop(stopCtx) cancel() if err != nil { r.signalFailure(err) return fmt.Errorf("stop task consumer before purge %q: %w", route.Queue, err) } delete(r.taskConsumers, route.Queue) } } 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) if r.sipUpdates != nil { select { case r.sipUpdates <- struct{}{}: default: // a pending notification has already requested a refresh } } return nil default: failure := fmt.Errorf("unexpected MQ control event %q", event.EventType) r.signalFailure(failure) return failure } } func (r *Runtime) 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 *Runtime) 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 *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error { r.consumersMu.Lock() defer r.consumersMu.Unlock() assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID) if err != nil { return err } wanted := make(map[string]store.AssignedTask, len(assigned)) for _, task := range assigned { route, err := tenant.TaskRoute(r.Bootstrap.DispatcherID, task.TaskID) if err != nil { return err } if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" || task.ControlState == "stopping" || task.ControlState == "stopped" || task.Status == "stopped" { continue } 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.Consumer 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 _ Publisher = (*mq.Broker)(nil)