package dispatcher import ( "context" "encoding/json" "errors" "fmt" "log/slog" "time" "git.ipao.vip/rogee/go-sip/internal/configread" "git.ipao.vip/rogee/go-sip/internal/contract" "git.ipao.vip/rogee/go-sip/internal/store" ) // ControlSpec is an Agent control instruction, not evidence that a // drain/hangup has completed. Never log a user-supplied reason or credentials. type ControlSpec struct { DispatcherID string TenantID int64 TaskID string Action string ActiveCallPolicy string Reason string } type ControlAgent interface { SendControl(context.Context, ControlSpec) error } type ControlController struct { DispatcherID string Store *store.Store Client *configread.Client Agent ControlAgent ApprovedSIP *configread.SIP ApprovedProviders *map[string]configread.Provider Now func() time.Time PurgeTaskQueue func(context.Context, string) (int, error) } // ProcessControl applies each delivered control independently: it has no // command_id, expected revision, or control-message deduplication. The state // barrier is durable before Agent dispatch; applied state and MQ outbox are // committed atomically after the Agent accepts the instruction. func (c *ControlController) ProcessControl(ctx context.Context, body []byte) error { if c == nil || c.DispatcherID == "" || c.Store == nil || c.Client == nil || c.Agent == nil || c.ApprovedSIP == nil || c.ApprovedProviders == nil || c.Now == nil { return errors.New("control processing requires Dispatcher, durable store, HTTP client, Agent, approved SIP/provider snapshots, and clock") } if err := contract.ValidateCurrent("mq", body); err != nil { return fmt.Errorf("invalid incoming task.control: %w", err) } var event struct { EventID string `json:"event_id"` EventType string `json:"event_type"` DispatcherID string `json:"dispatcher_id"` TenantID int64 `json:"tenant_id"` IssuedAt string `json:"issued_at"` Payload struct { TaskID string `json:"task_id"` Action string `json:"action"` Reason string `json:"reason"` Options struct { ActiveCallPolicy string `json:"active_call_policy"` } `json:"options"` } `json:"payload"` } if err := json.Unmarshal(body, &event); err != nil { return fmt.Errorf("decode task.control: %w", err) } if event.EventType != "task.control" || event.DispatcherID != c.DispatcherID || event.Payload.TaskID == "" || event.Payload.Action == "" { return errors.New("control event type, task, or Dispatcher owner mismatch") } issuedAt, err := time.Parse(time.RFC3339Nano, event.IssuedAt) if err != nil { return fmt.Errorf("invalid task.control issued_at: %w", err) } if issuedAt.After(c.Now()) { return fmt.Errorf("control event %q issued_at is in the future; do not apply early", event.EventID) } policy := event.Payload.Options.ActiveCallPolicy if event.Payload.Action == "start" || event.Payload.Action == "resume" { if policy != "" { return errors.New("start/resume control must not carry an active-call policy") } // Task reads use the boot-approved global snapshots; the durable // SIP and pause/stop barriers remain authoritative for admission. snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID, *c.ApprovedSIP, *c.ApprovedProviders) if err != nil { return fmt.Errorf("fresh %s task configuration: %w", event.Payload.Action, err) } if snapshot.Task.Status != "running" { return fmt.Errorf("fresh %s task is not running", event.Payload.Action) } if err := validateAISnapshot(snapshot); err != nil { return err } if err := c.Store.SaveSnapshot(snapshot); err != nil { return fmt.Errorf("bind fresh %s task: %w", event.Payload.Action, err) } if event.Payload.Action == "start" { if err := c.Store.CompleteStart(snapshot, event.EventID); err != nil { if errors.Is(err, store.ErrControlRejected) { return c.Store.RejectControl(event.DispatcherID, event.TenantID, event.EventID) } return fmt.Errorf("start task %q: %w", event.Payload.TaskID, err) } return nil } if err := c.Store.ApplyDiscoveryPage(event.DispatcherID, []configread.DiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: snapshot.Task.Status}}); err != nil { return fmt.Errorf("bind fresh resume task revision: %w", err) } } else { if policy == "" { policy = "hangup" } if policy != "hangup" && policy != "drain" { return errors.New("unapproved active-call policy") } } if err := c.Store.PrepareControl(event.DispatcherID, event.TenantID, event.Payload.TaskID, event.Payload.Action); err != nil { if errors.Is(err, store.ErrControlRejected) { if rejectErr := c.Store.RejectControl(event.DispatcherID, event.TenantID, event.EventID); rejectErr != nil { return fmt.Errorf("persist rejected control %q: %w", event.EventID, rejectErr) } return nil } return fmt.Errorf("prepare task control %q: %w", event.EventID, err) } if event.Payload.Action == "stop" { if c.PurgeTaskQueue == nil { return errors.New("stop requires a task queue purger") } count, err := c.PurgeTaskQueue(ctx, event.Payload.TaskID) if err != nil { return fmt.Errorf("purge stopped task %q backlog: %w", event.Payload.TaskID, err) } // A purge covers Ready messages only. The durable stop barrier also // rejects any later or already accepted deliveries. slog.Info("stopped task backlog purged", "dispatcher_id", event.DispatcherID, "task_id", event.Payload.TaskID, "ready_messages", count) } spec := ControlSpec{ DispatcherID: event.DispatcherID, TenantID: event.TenantID, TaskID: event.Payload.TaskID, Action: event.Payload.Action, ActiveCallPolicy: policy, Reason: event.Payload.Reason, } if err := c.Agent.SendControl(ctx, spec); err != nil { return fmt.Errorf("Agent control %q dispatch failed: %w", event.EventID, err) } if err := c.Store.CompleteControl(event.DispatcherID, event.TenantID, event.Payload.TaskID, event.Payload.Action, event.EventID); err != nil { return fmt.Errorf("Agent control %q dispatched but acknowledgment not durable: %w", event.EventID, err) } return nil }