Files
go-sip/internal/dispatcher/control.go
T

155 lines
6.1 KiB
Go

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
}