Use event-driven SIP reload and approved task bindings
This commit is contained in:
@@ -138,9 +138,9 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
|
||||
DispatcherID: settings.DispatcherID, Store: database, Originator: originator, Publisher: broker, Now: time.Now,
|
||||
},
|
||||
Control: dispatcher.ControlController{
|
||||
DispatcherID: settings.DispatcherID, Store: database, Client: reader, Agent: originator, VerifySIP: originator.VerifySIP, Now: time.Now,
|
||||
DispatcherID: settings.DispatcherID, Store: database, Client: reader, Agent: originator, Now: time.Now,
|
||||
},
|
||||
PollInterval: time.Second, SIPInterval: time.Minute, Logger: slog.Default(),
|
||||
PollInterval: time.Second, Logger: slog.Default(),
|
||||
}
|
||||
listener, err := net.Listen("tcp", settings.Listen)
|
||||
if err != nil {
|
||||
|
||||
@@ -285,7 +285,7 @@ func TestDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) {
|
||||
finished := make(chan error, 1)
|
||||
go func() { finished <- command.Execute() }()
|
||||
deadline := time.After(8 * time.Second)
|
||||
for sipReads.Load() < 2 || discoveryReads.Load() < 2 || taskReads.Load() == 0 || providerReads.Load() == 0 || quotaReads.Load() == 0 {
|
||||
for sipReads.Load() < 1 || discoveryReads.Load() < 2 || taskReads.Load() == 0 || providerReads.Load() == 0 || quotaReads.Load() == 0 {
|
||||
select {
|
||||
case err := <-finished:
|
||||
cancel()
|
||||
|
||||
@@ -119,18 +119,16 @@ func (c *Client) ReadSIP(ctx context.Context) (SIP, error) {
|
||||
return sip, nil
|
||||
}
|
||||
|
||||
// ReadTask reads only the agreed read-only HTTP resources. It never
|
||||
// substitutes stale data or the historical MQ configuration contract.
|
||||
func (c *Client) ReadTask(ctx context.Context, taskID string, tenantID int64) (Snapshot, error) {
|
||||
// ReadTask reads task, provider and quota resources using the already approved
|
||||
// SIP snapshot. Only bootstrap and sip.config may fetch the full SIP resource.
|
||||
func (c *Client) ReadTask(ctx context.Context, taskID string, tenantID int64, sip SIP) (Snapshot, error) {
|
||||
if taskID == "" || tenantID <= 0 {
|
||||
return Snapshot{}, errors.New("task ID and positive tenant ID are required")
|
||||
}
|
||||
var result Snapshot
|
||||
var err error
|
||||
result.SIP, err = c.ReadSIP(ctx)
|
||||
if err != nil {
|
||||
return Snapshot{}, err
|
||||
if sip.DispatcherID != c.dispatcherID || sip.Revision <= 0 {
|
||||
return Snapshot{}, errors.New("task requires the current approved SIP snapshot")
|
||||
}
|
||||
result := Snapshot{SIP: sip}
|
||||
var providers providerResponse
|
||||
if err := c.readResource(ctx, configReadPath+"/ai-providers", "ai_providers", &providers); err != nil {
|
||||
return Snapshot{}, fmt.Errorf("read AI providers: %w", err)
|
||||
|
||||
@@ -22,6 +22,15 @@ func example(t *testing.T, name string) []byte {
|
||||
return raw
|
||||
}
|
||||
|
||||
func approvedSIP(t *testing.T) SIP {
|
||||
t.Helper()
|
||||
var sip SIP
|
||||
if err := json.Unmarshal(example(t, "config-read-sip"), &sip); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return sip
|
||||
}
|
||||
|
||||
func TestReadTaskSnapshot(t *testing.T) {
|
||||
for _, mode := range []string{"asr", "full"} {
|
||||
t.Run(mode, func(t *testing.T) {
|
||||
@@ -53,11 +62,11 @@ func TestReadTaskSnapshot(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
snapshot, err := client.ReadTask(context.Background(), "task-"+mode, 1001)
|
||||
snapshot, err := client.ReadTask(context.Background(), "task-"+mode, 1001, approvedSIP(t))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if calls != 4 || snapshot.SIP.Revision != 8 || snapshot.Task.TenantID != 1001 || snapshot.Quota.TenantID != 1001 {
|
||||
if calls != 3 || snapshot.SIP.Revision != 8 || snapshot.Task.TenantID != 1001 || snapshot.Quota.TenantID != 1001 {
|
||||
t.Fatalf("unexpected request count or identity: calls=%d snapshot=%+v", calls, snapshot)
|
||||
}
|
||||
if got := snapshot.Providers["asr-example"].Credential; got != "example-only-not-a-real-secret" {
|
||||
@@ -114,7 +123,6 @@ func TestReadTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) {
|
||||
"/internal/v1/dispatcher/tenant/1001/quota": example(t, "config-read-quota"),
|
||||
}
|
||||
for _, tc := range []struct{ path, owner string }{
|
||||
{"/internal/v1/dispatcher/sip", "SIP"},
|
||||
{"/internal/v1/dispatcher/ai-providers", "AI providers"},
|
||||
{"/internal/v1/dispatcher/task/task-asr", "task configuration"},
|
||||
{"/internal/v1/dispatcher/tenant/1001/quota", "tenant quota"},
|
||||
@@ -140,11 +148,11 @@ func TestReadTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := client.ReadTask(context.Background(), "task-asr", 1001); err != nil {
|
||||
if _, err := client.ReadTask(context.Background(), "task-asr", 1001, approvedSIP(t)); err != nil {
|
||||
t.Fatalf("approved resources failed before the outage: %v", err)
|
||||
}
|
||||
unavailable.Store(true)
|
||||
_, err = client.ReadTask(context.Background(), "task-asr", 1001)
|
||||
_, err = client.ReadTask(context.Background(), "task-asr", 1001, approvedSIP(t))
|
||||
var httpErr *HTTPError
|
||||
if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable {
|
||||
t.Fatalf("unavailable %s reused cached configuration or lost status: %v", tc.owner, err)
|
||||
@@ -157,7 +165,7 @@ func TestReadTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestReadTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) {
|
||||
for _, mode := range []string{"owner", "sip-revision", "provider-ref", "provider-disabled", "provider-wrong-role"} {
|
||||
for _, mode := range []string{"owner", "provider-ref", "provider-disabled", "provider-wrong-role"} {
|
||||
t.Run(mode, func(t *testing.T) {
|
||||
responses := map[string][]byte{
|
||||
"/internal/v1/dispatcher/sip": example(t, "config-read-sip"),
|
||||
@@ -168,8 +176,6 @@ func TestReadTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) {
|
||||
switch mode {
|
||||
case "owner":
|
||||
responses["/internal/v1/dispatcher/task/task-asr"] = []byte(strings.Replace(string(responses["/internal/v1/dispatcher/task/task-asr"]), `"tenant_id":1001`, `"tenant_id":1002`, 1))
|
||||
case "sip-revision":
|
||||
responses["/internal/v1/dispatcher/sip"] = example(t, "invalid/config-read-sip-missing-revision")
|
||||
case "provider-ref":
|
||||
responses["/internal/v1/dispatcher/ai-providers"] = example(t, "invalid/config-read-provider-ref")
|
||||
case "provider-disabled":
|
||||
@@ -186,7 +192,7 @@ func TestReadTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := client.ReadTask(context.Background(), "task-asr", 1001); err == nil || strings.Contains(err.Error(), "example-only-not-a-real-secret") {
|
||||
if _, err := client.ReadTask(context.Background(), "task-asr", 1001, approvedSIP(t)); err == nil || strings.Contains(err.Error(), "example-only-not-a-real-secret") {
|
||||
t.Fatalf("expected redacted fail-closed error, got %v", err)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -2,10 +2,8 @@ package dispatcher
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"reflect"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
@@ -58,19 +56,12 @@ func (b Bootstrap) Run(ctx context.Context) error {
|
||||
if task.Status == "stopped" {
|
||||
continue
|
||||
}
|
||||
snapshot, err := b.Client.ReadTask(ctx, task.TaskID, task.TenantID)
|
||||
snapshot, err := b.Client.ReadTask(ctx, task.TaskID, task.TenantID, sip)
|
||||
if err != nil {
|
||||
return fmt.Errorf("read assigned task %q: %w", task.TaskID, err)
|
||||
}
|
||||
var approvedTrunks, taskTrunks any
|
||||
if err := json.Unmarshal(sip.Trunks, &approvedTrunks); err != nil {
|
||||
return fmt.Errorf("decode approved SIP trunks: %w", err)
|
||||
}
|
||||
if err := json.Unmarshal(snapshot.SIP.Trunks, &taskTrunks); err != nil {
|
||||
return fmt.Errorf("decode task SIP trunks: %w", err)
|
||||
}
|
||||
if snapshot.SIP.Revision != sip.Revision || !reflect.DeepEqual(taskTrunks, approvedTrunks) || snapshot.Task.TaskRevision != task.TaskRevision || snapshot.Task.Status != task.Status {
|
||||
return fmt.Errorf("task %q configuration differs from approved SIP/discovery snapshot", task.TaskID)
|
||||
if snapshot.Task.TaskRevision != task.TaskRevision || snapshot.Task.Status != task.Status {
|
||||
return fmt.Errorf("task %q configuration differs from approved discovery snapshot", task.TaskID)
|
||||
}
|
||||
if err := validateAISnapshot(snapshot); err != nil {
|
||||
return err
|
||||
@@ -79,6 +70,9 @@ func (b Bootstrap) Run(ctx context.Context) error {
|
||||
return fmt.Errorf("persist immutable task %q: %w", task.TaskID, err)
|
||||
}
|
||||
}
|
||||
if b.SIP != nil {
|
||||
*b.SIP = sip
|
||||
}
|
||||
if err := b.DrainControls(ctx); err != nil {
|
||||
return fmt.Errorf("drain assigned Dispatcher control queue: %w", err)
|
||||
}
|
||||
@@ -88,8 +82,14 @@ func (b Bootstrap) Run(ctx context.Context) error {
|
||||
if b.Cursor != nil {
|
||||
*b.Cursor = cursor
|
||||
}
|
||||
if b.SIP != nil {
|
||||
*b.SIP = sip
|
||||
// A sip.config delivered during control drain remains fenced and is
|
||||
// reconciled by the runtime notification loop rather than aborting boot.
|
||||
_, pending, err := b.Store.SIPState(b.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if pending != 0 {
|
||||
return nil
|
||||
}
|
||||
if err := b.Store.MarkReadyForSIP(b.DispatcherID, sip.Revision); err != nil {
|
||||
return fmt.Errorf("open admitted tasks: %w", err)
|
||||
|
||||
@@ -88,11 +88,9 @@ func TestBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) {
|
||||
if err != nil || !admitted || !verifierCalled || !drained || cursor != "opaque-end-token" || approvedSIP.Revision != 8 || len(called) < 6 || called[0] != "/internal/v1/dispatcher/sip" {
|
||||
t.Fatalf("bootstrap: admitted=%v, err=%v, verified=%v, drained=%v, paths=%v", admitted, err, verifierCalled, drained, called)
|
||||
}
|
||||
if err := db.NoteSIPChange(id, 9); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := bootstrap.Run(context.Background()); !errors.Is(err, store.ErrSIPPending) {
|
||||
t.Fatalf("startup lost durable newer SIP notification: %v", err)
|
||||
bootstrap.DrainControls = func(context.Context) error { return db.NoteSIPChange(id, 9) }
|
||||
if err := bootstrap.Run(context.Background()); err != nil {
|
||||
t.Fatalf("notification delivered during startup must remain pending for runtime: %v", err)
|
||||
}
|
||||
if cursor != "opaque-end-token" || approvedSIP.Revision != 8 {
|
||||
t.Fatalf("pending SIP startup lost verified snapshot/cursor: revision=%d cursor=%q", approvedSIP.Revision, cursor)
|
||||
|
||||
@@ -32,7 +32,7 @@ type ControlController struct {
|
||||
Store *store.Store
|
||||
Client *configread.Client
|
||||
Agent ControlAgent
|
||||
VerifySIP func(context.Context, configread.SIP) error
|
||||
ApprovedSIP *configread.SIP
|
||||
Now func() time.Time
|
||||
}
|
||||
|
||||
@@ -41,8 +41,8 @@ type ControlController struct {
|
||||
// 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.VerifySIP == nil || c.Now == nil {
|
||||
return errors.New("control processing requires Dispatcher, durable store, HTTP client, Agent, SIP verifier, and clock")
|
||||
if c == nil || c.DispatcherID == "" || c.Store == nil || c.Client == nil || c.Agent == nil || c.ApprovedSIP == nil || c.Now == nil {
|
||||
return errors.New("control processing requires Dispatcher, durable store, HTTP client, Agent, approved SIP snapshot, and clock")
|
||||
}
|
||||
if err := contract.ValidateCurrent("mq", body); err != nil {
|
||||
return fmt.Errorf("invalid incoming task.control: %w", err)
|
||||
@@ -80,18 +80,15 @@ func (c *ControlController) ProcessControl(ctx context.Context, body []byte) err
|
||||
if policy != "" {
|
||||
return errors.New("start/resume control must not carry an active-call policy")
|
||||
}
|
||||
// Only a fresh approved task plus applied SIP snapshot may authorize
|
||||
// new admission; the durable pause/stop barrier remains authoritative.
|
||||
snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID)
|
||||
// Task reads use the already approved SIP snapshot; the durable
|
||||
// SIP and pause/stop barriers remain authoritative for admission.
|
||||
snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID, *c.ApprovedSIP)
|
||||
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 := c.VerifySIP(ctx, snapshot.SIP); err != nil {
|
||||
return fmt.Errorf("%s SIP revision not applied: %w", event.Payload.Action, err)
|
||||
}
|
||||
if err := validateAISnapshot(snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -78,7 +78,7 @@ func newControlFixture(t *testing.T) (*ControlController, *fakeControlAgent, *st
|
||||
t.Fatal(err)
|
||||
}
|
||||
agent := &fakeControlAgent{}
|
||||
controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, VerifySIP: func(context.Context, configread.SIP) error { return nil }, Now: func() time.Time { return monday(9, 30) }}
|
||||
controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, ApprovedSIP: &snapshot.SIP, Now: func() time.Time { return monday(9, 30) }}
|
||||
return controller, agent, s
|
||||
}
|
||||
|
||||
@@ -179,33 +179,22 @@ func TestStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox(t *test
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) {
|
||||
func TestResumeUsesApprovedSIPButPendingRevisionStillClosesAdmission(t *testing.T) {
|
||||
controller, agent, s := newControlFixture(t)
|
||||
if err := controller.ProcessControl(context.Background(), controlBody(t, "pause-before-resume", "pause", "drain")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
approved := controller.VerifySIP
|
||||
controller.VerifySIP = func(context.Context, configread.SIP) error {
|
||||
return errors.New("SIP revision not yet applied by Asterisk")
|
||||
if err := s.NoteSIPChange(controller.DispatcherID, 9); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := controller.ProcessControl(context.Background(), controlBody(t, "resume-after-load", "resume", "")); err == nil {
|
||||
t.Fatal("resumed task without applied SIP")
|
||||
}
|
||||
if len(agent.calls) != 1 {
|
||||
t.Fatalf("unverified resume reached Agent: %+v", agent.calls)
|
||||
}
|
||||
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted {
|
||||
t.Fatalf("failed SIP check reopened admission: %v %v", admitted, err)
|
||||
}
|
||||
controller.VerifySIP = approved
|
||||
if err := controller.ProcessControl(context.Background(), controlBody(t, "resume-after-load", "resume", "")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(agent.calls) != 2 || agent.calls[1].Action != "resume" || agent.calls[1].ActiveCallPolicy != "" {
|
||||
t.Fatalf("valid fresh resume did not dispatch: %+v", agent.calls)
|
||||
if len(agent.calls) != 2 || agent.calls[1].Action != "resume" {
|
||||
t.Fatalf("resume did not reach Agent: %+v", agent.calls)
|
||||
}
|
||||
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || !admitted {
|
||||
t.Fatalf("verified resume did not reopen admission: %v %v", admitted, err)
|
||||
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted {
|
||||
t.Fatalf("pending SIP revision reopened admission: %v %v", admitted, err)
|
||||
}
|
||||
outbox, err := s.ListPendingOutbox(controller.DispatcherID)
|
||||
if err != nil || len(outbox) != 2 || outbox[1].EventID != "resume-after-load" {
|
||||
|
||||
@@ -23,20 +23,20 @@ type Runtime struct {
|
||||
Execute ExecuteController
|
||||
Control ControlController
|
||||
PollInterval time.Duration
|
||||
SIPInterval 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
|
||||
gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction
|
||||
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, task queues, or the SQLite file.
|
||||
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.SIPInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil {
|
||||
return errors.New("runtime requires one Dispatcher, durable state, verified SIP, independent clocks, and configured polling; external control drain is forbidden")
|
||||
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
|
||||
@@ -45,6 +45,7 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
|
||||
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)
|
||||
var controlConsumer *mq.Consumer
|
||||
@@ -68,6 +69,7 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
|
||||
|
||||
var sip configread.SIP
|
||||
r.Bootstrap.SIP = &sip
|
||||
r.Control.ApprovedSIP = &sip
|
||||
r.Bootstrap.DrainControls = func(ctx context.Context) error {
|
||||
_, err := r.Broker.DrainControlPredeclared(ctx, r.Broker.ControlQueue(), r.handleControl)
|
||||
if err != nil {
|
||||
@@ -99,8 +101,6 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
|
||||
}
|
||||
poll := time.NewTicker(r.PollInterval)
|
||||
defer poll.Stop()
|
||||
sipTick := time.NewTicker(r.SIPInterval)
|
||||
defer sipTick.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -108,6 +108,9 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
|
||||
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
|
||||
}
|
||||
@@ -117,12 +120,9 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
|
||||
if err := r.syncTaskConsumers(ctx, consumers); err != nil {
|
||||
return err
|
||||
}
|
||||
case <-sipTick.C:
|
||||
case <-r.sipUpdates:
|
||||
if err := r.refreshSIP(ctx, &sip); err != nil {
|
||||
return fmt.Errorf("refresh approved SIP: %w", err)
|
||||
}
|
||||
if err := r.syncTaskConsumers(ctx, consumers); err != nil {
|
||||
return err
|
||||
return fmt.Errorf("apply SIP notification: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -198,6 +198,12 @@ func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) erro
|
||||
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)
|
||||
|
||||
@@ -100,12 +100,13 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var taskListReads atomic.Int32
|
||||
var taskListReads, sipReads atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
var body []byte
|
||||
switch r.URL.Path {
|
||||
case "/internal/v1/dispatcher/sip":
|
||||
sipReads.Add(1)
|
||||
body = sipJSON
|
||||
case "/internal/v1/dispatcher/tasks":
|
||||
taskListReads.Add(1)
|
||||
@@ -165,8 +166,8 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
|
||||
}
|
||||
return monday(8, 59)
|
||||
}},
|
||||
Control: ControlController{DispatcherID: id, Store: db, Client: client, Agent: agent, VerifySIP: verify, Now: func() time.Time { return monday(9, 30) }},
|
||||
PollInterval: 30 * time.Millisecond, SIPInterval: 120 * time.Millisecond, Logger: slog.Default(),
|
||||
Control: ControlController{DispatcherID: id, Store: db, Client: client, Agent: agent, Now: func() time.Time { return monday(9, 30) }},
|
||||
PollInterval: 30 * time.Millisecond, Logger: slog.Default(),
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
finished := make(chan struct{})
|
||||
@@ -285,6 +286,9 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
|
||||
if got := taskListReads.Load(); got != 2 {
|
||||
t.Fatalf("unexpected periodic task-list reads: %d", got)
|
||||
}
|
||||
if got := sipReads.Load(); got != 1 {
|
||||
t.Fatalf("task start/resume or idle polling reread full SIP: %d", got)
|
||||
}
|
||||
// A temporary task rule wait is retained durably, then its consumer
|
||||
// stops so further instructions remain in SaaS's task queue. Admission
|
||||
// resumes from the original identity when the configured window opens.
|
||||
|
||||
@@ -56,8 +56,11 @@ func TestBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) {
|
||||
VerifySIP: func(context.Context, configread.SIP) error { return nil },
|
||||
DrainControls: func(context.Context) error { return nil },
|
||||
}
|
||||
if err := bootstrap.Run(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := bootstrap.Run(context.Background()); err == nil {
|
||||
t.Fatal("SIP content changed without revision, but admission opened")
|
||||
t.Fatal("SIP content changed without revision on restart, but admission opened")
|
||||
}
|
||||
if admitted, err := db.CanAdmit(id, 1001, "task-asr"); err != nil || admitted {
|
||||
t.Fatalf("admission opened after SIP content conflict: %v, %v", admitted, err)
|
||||
|
||||
@@ -10,13 +10,27 @@ import (
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
)
|
||||
|
||||
// refreshSIP keeps control/result delivery running while a new approved SIP
|
||||
// revision waits for old calls to drain and for Agent/Asterisk to load it.
|
||||
// No HTTP snapshot or notification by itself authorizes a real call.
|
||||
// refreshSIP processes only a durable sip.config notification. A pending
|
||||
// revision stays fenced until existing calls drain and Agent confirms loading.
|
||||
func (r *Runtime) refreshSIP(ctx context.Context, approvedSIP *configread.SIP) error {
|
||||
if r == nil || approvedSIP == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Logger == nil {
|
||||
return errors.New("SIP refresh requires durable state and applied-revision verifier")
|
||||
}
|
||||
_, pending, err := r.Bootstrap.Store.SIPState(r.Bootstrap.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if pending == 0 {
|
||||
return nil
|
||||
}
|
||||
occupied, err := r.Bootstrap.Store.OccupiedCalls(r.Bootstrap.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if occupied != 0 {
|
||||
r.Logger.Debug("SIP reload waits for confirmed old-call drain", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "occupied_calls", occupied)
|
||||
return nil
|
||||
}
|
||||
current, err := r.Bootstrap.Client.ReadSIP(ctx)
|
||||
if err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("read approved SIP full snapshot: %w", err))
|
||||
@@ -36,72 +50,16 @@ func (r *Runtime) refreshSIP(ctx context.Context, approvedSIP *configread.SIP) e
|
||||
if current.Revision == old.Revision && !bytes.Equal(oldJSON, currentJSON) {
|
||||
return r.closeSIPAdmission(fmt.Errorf("approved SIP revision %d changed content without a new revision", current.Revision))
|
||||
}
|
||||
if current.Revision > old.Revision {
|
||||
r.gate.Lock()
|
||||
err := r.Bootstrap.Store.NoteSIPChange(r.Bootstrap.DispatcherID, current.Revision)
|
||||
r.gate.Unlock()
|
||||
if err != nil {
|
||||
return fmt.Errorf("fence newly discovered SIP revision: %w", err)
|
||||
}
|
||||
}
|
||||
_, pending, err := r.Bootstrap.Store.SIPState(r.Bootstrap.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if pending == 0 {
|
||||
return nil
|
||||
}
|
||||
if current.Revision < pending {
|
||||
r.Logger.Warn("SaaS SIP snapshot is behind durable SIP notification; admission stays closed", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "available_revision", current.Revision)
|
||||
return nil
|
||||
}
|
||||
occupied, err := r.Bootstrap.Store.OccupiedCalls(r.Bootstrap.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if occupied != 0 {
|
||||
r.Logger.Info("SIP reload waits for confirmed old-call drain", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "occupied_calls", occupied)
|
||||
r.Logger.Debug("SaaS SIP snapshot is behind durable SIP notification; admission stays closed", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "available_revision", current.Revision)
|
||||
return nil
|
||||
}
|
||||
if err := r.Bootstrap.VerifySIP(ctx, current); err != nil {
|
||||
r.Logger.Warn("SIP reload waits for actual Agent/Asterisk revision", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "error", err)
|
||||
r.Logger.Debug("SIP reload waits for actual Agent/Asterisk revision", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "error", err)
|
||||
return nil
|
||||
}
|
||||
// A full discovery snapshot is required after an approved SIP revision
|
||||
// change; a delta page cannot prove that all assigned tasks use the same
|
||||
// loaded revision. MQ control consumption remains active during this read.
|
||||
tasks, _, err := r.Bootstrap.Client.ReadAllTasks(ctx)
|
||||
if err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("reload complete assigned task list for SIP: %w", err))
|
||||
}
|
||||
if err := r.Bootstrap.Store.ApplyDiscoverySnapshot(r.Bootstrap.DispatcherID, tasks); err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("persist complete assigned task list for SIP: %w", err))
|
||||
}
|
||||
for _, task := range tasks {
|
||||
if task.Status == "stopped" {
|
||||
continue
|
||||
}
|
||||
snapshot, err := r.Bootstrap.Client.ReadTask(ctx, task.TaskID, task.TenantID)
|
||||
if err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("reload task %q after SIP revision: %w", task.TaskID, err))
|
||||
}
|
||||
if snapshot.Task.TaskRevision != task.TaskRevision || snapshot.Task.Status != task.Status {
|
||||
return r.closeSIPAdmission(fmt.Errorf("task %q discovery conflicts with SIP reload configuration", task.TaskID))
|
||||
}
|
||||
actualJSON, err := json.Marshal(snapshot.SIP)
|
||||
if err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("encode task %q SIP binding: %w", task.TaskID, err))
|
||||
}
|
||||
if !bytes.Equal(currentJSON, actualJSON) {
|
||||
return r.closeSIPAdmission(fmt.Errorf("task %q is not bound to approved SIP revision %d", task.TaskID, current.Revision))
|
||||
}
|
||||
if err := validateAISnapshot(snapshot); err != nil {
|
||||
return r.closeSIPAdmission(err)
|
||||
}
|
||||
if err := r.Bootstrap.Store.SaveSnapshot(snapshot); err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("persist task %q SIP binding: %w", task.TaskID, err))
|
||||
}
|
||||
}
|
||||
r.gate.Lock()
|
||||
defer r.gate.Unlock()
|
||||
_, latestPending, err := r.Bootstrap.Store.SIPState(r.Bootstrap.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -110,8 +68,10 @@ func (r *Runtime) refreshSIP(ctx context.Context, approvedSIP *configread.SIP) e
|
||||
r.Logger.Info("newer SIP notification arrived during reload; admission remains closed", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", latestPending)
|
||||
return nil
|
||||
}
|
||||
// MarkReadyForSIP rechecks the drain and every durable task binding in
|
||||
// one SQLite transaction. No old version can reopen admission here.
|
||||
if err := r.Bootstrap.Store.RebindSIP(r.Bootstrap.DispatcherID, current); err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("bind approved SIP to existing task snapshots: %w", err))
|
||||
}
|
||||
// This transaction rechecks drain, pending revision and every binding.
|
||||
if err := r.Bootstrap.Store.MarkReadyForSIP(r.Bootstrap.DispatcherID, current.Revision); err != nil {
|
||||
return r.closeSIPAdmission(fmt.Errorf("verify and commit reloaded SIP revision: %w", err))
|
||||
}
|
||||
|
||||
@@ -29,18 +29,6 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing.
|
||||
switch r.URL.Path {
|
||||
case "/internal/v1/dispatcher/sip":
|
||||
body = newSIP
|
||||
case "/internal/v1/dispatcher/tasks":
|
||||
if r.URL.Query().Get("after") == "" {
|
||||
body = configExample(t, "task-discovery-page")
|
||||
} else {
|
||||
body = configExample(t, "task-discovery-end")
|
||||
}
|
||||
case "/internal/v1/dispatcher/ai-providers":
|
||||
body = configExample(t, "config-read-providers")
|
||||
case "/internal/v1/dispatcher/task/task-asr":
|
||||
body = configExample(t, "config-read-task-asr")
|
||||
case "/internal/v1/dispatcher/tenant/1001/quota":
|
||||
body = configExample(t, "config-read-quota")
|
||||
default:
|
||||
t.Errorf("unexpected HTTP path %s", r.URL.Path)
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
@@ -88,4 +76,8 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing.
|
||||
if approved.Revision != 9 {
|
||||
t.Fatalf("reloaded SIP not bound: rev=%d", approved.Revision)
|
||||
}
|
||||
bound, err := s.ReadSnapshot(executor.DispatcherID, 1001, "task-asr")
|
||||
if err != nil || bound.SIP.Revision != 9 {
|
||||
t.Fatalf("existing task did not inherit the approved SIP revision: %+v %v", bound.SIP, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,8 +2,12 @@ package store
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"reflect"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
)
|
||||
|
||||
var ErrSIPPending = errors.New("new SIP revision is pending drain or verified load")
|
||||
@@ -59,6 +63,84 @@ func (s *Store) OccupiedCalls(dispatcherID string) (int64, error) {
|
||||
return count, nil
|
||||
}
|
||||
|
||||
// RebindSIP updates the already approved SIP data in every persisted task
|
||||
// snapshot without refetching or altering immutable task/provider content.
|
||||
func (s *Store) RebindSIP(dispatcherID string, sip configread.SIP) error {
|
||||
if dispatcherID == "" || sip.DispatcherID != dispatcherID || sip.Revision <= 0 {
|
||||
return errors.New("invalid approved SIP binding")
|
||||
}
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
rows, err := tx.Query(`SELECT tenant_id,task_id,sip_revision,snapshot_json FROM dispatcher_configs WHERE dispatcher_id=?`, dispatcherID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("read task SIP bindings: %w", err)
|
||||
}
|
||||
type binding struct {
|
||||
tenantID int64
|
||||
taskID string
|
||||
body []byte
|
||||
}
|
||||
var bindings []binding
|
||||
for rows.Next() {
|
||||
var b binding
|
||||
var revision int64
|
||||
if err := rows.Scan(&b.tenantID, &b.taskID, &revision, &b.body); err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("scan task SIP binding: %w", err)
|
||||
}
|
||||
if revision > sip.Revision {
|
||||
rows.Close()
|
||||
return fmt.Errorf("task %q SIP revision regressed", b.taskID)
|
||||
}
|
||||
var data struct {
|
||||
SIP configread.SIP `json:"sip"`
|
||||
}
|
||||
if err := json.Unmarshal(b.body, &data); err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("decode task %q SIP binding: %w", b.taskID, err)
|
||||
}
|
||||
if revision == sip.Revision && !reflect.DeepEqual(data.SIP, sip) {
|
||||
rows.Close()
|
||||
return fmt.Errorf("task %q SIP revision %d changed content", b.taskID, revision)
|
||||
}
|
||||
var fields map[string]json.RawMessage
|
||||
if err := json.Unmarshal(b.body, &fields); err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("decode task %q snapshot: %w", b.taskID, err)
|
||||
}
|
||||
fields["sip"], err = json.Marshal(sip)
|
||||
if err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("encode approved SIP: %w", err)
|
||||
}
|
||||
b.body, err = json.Marshal(fields)
|
||||
if err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("encode task %q SIP binding: %w", b.taskID, err)
|
||||
}
|
||||
bindings = append(bindings, b)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("iterate task SIP bindings: %w", err)
|
||||
}
|
||||
if err := rows.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, b := range bindings {
|
||||
if _, err := tx.Exec(`UPDATE dispatcher_configs SET sip_revision=?,snapshot_json=? WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, sip.Revision, b.body, dispatcherID, b.tenantID, b.taskID); err != nil {
|
||||
return fmt.Errorf("rebind task %q SIP revision: %w", b.taskID, err)
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("commit approved SIP bindings: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// MarkReadyForSIP opens admission only after all present task bindings match
|
||||
// this exact verified SIP revision and any prior SIP calls have drained.
|
||||
// The loaded Agent/Asterisk state is independently checked by the caller.
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
@@ -68,6 +69,34 @@ func TestSIPNotificationRequiresFullDrainAndExactLoadedRevision(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRebindSIPPreservesTaskAndRejectsSameRevisionConflict(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
original, err := s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updated := original.SIP
|
||||
updated.Revision++
|
||||
if err := s.NoteSIPChange(currentDispatcherID, updated.Revision); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.RebindSIP(currentDispatcherID, updated); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
loaded, err := s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil || loaded.SIP.Revision != updated.Revision || !bytes.Equal(loaded.Task.Raw, original.Task.Raw) || loaded.Quota != original.Quota {
|
||||
t.Fatalf("SIP update changed task or quota: %+v %v", loaded, err)
|
||||
}
|
||||
if err := s.MarkReadyForSIP(currentDispatcherID, updated.Revision); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
altered := updated
|
||||
altered.Trunks = []byte(`[]`)
|
||||
if err := s.RebindSIP(currentDispatcherID, altered); err == nil {
|
||||
t.Fatal("same SIP revision changed content")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPendingSIPChangeSurvivesRestartWithoutOldDataDeletion(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "state.db")
|
||||
s, err := Open(path)
|
||||
|
||||
Reference in New Issue
Block a user