From 9fe539569ec4ee061b45dedd7266fcb2405959db Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 1 Oct 2026 13:52:48 +0800 Subject: [PATCH] Use event-driven SIP reload and approved task bindings --- cmd/sip-go-agent/dispatcher_command.go | 4 +- .../dispatcher_integration_test.go | 2 +- internal/configread/snapshots.go | 14 ++- internal/configread/snapshots_test.go | 24 +++-- internal/dispatcher/config.go | 28 +++--- internal/dispatcher/config_test.go | 8 +- internal/dispatcher/control.go | 15 ++-- internal/dispatcher/control_test.go | 27 ++---- internal/dispatcher/runtime.go | 34 ++++--- .../dispatcher/runtime_integration_test.go | 10 ++- internal/dispatcher/sip_conflict_test.go | 5 +- internal/dispatcher/sip_reload.go | 90 ++++++------------- internal/dispatcher/sip_runtime_test.go | 16 +--- internal/store/sip.go | 82 +++++++++++++++++ internal/store/sip_change_test.go | 29 ++++++ 15 files changed, 226 insertions(+), 162 deletions(-) diff --git a/cmd/sip-go-agent/dispatcher_command.go b/cmd/sip-go-agent/dispatcher_command.go index 295de50..38d333b 100644 --- a/cmd/sip-go-agent/dispatcher_command.go +++ b/cmd/sip-go-agent/dispatcher_command.go @@ -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 { diff --git a/cmd/sip-go-agent/dispatcher_integration_test.go b/cmd/sip-go-agent/dispatcher_integration_test.go index ec9dd49..58afe21 100644 --- a/cmd/sip-go-agent/dispatcher_integration_test.go +++ b/cmd/sip-go-agent/dispatcher_integration_test.go @@ -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() diff --git a/internal/configread/snapshots.go b/internal/configread/snapshots.go index da7f0d6..60b99ca 100644 --- a/internal/configread/snapshots.go +++ b/internal/configread/snapshots.go @@ -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) diff --git a/internal/configread/snapshots_test.go b/internal/configread/snapshots_test.go index 8267b64..f1a0e51 100644 --- a/internal/configread/snapshots_test.go +++ b/internal/configread/snapshots_test.go @@ -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) } }) diff --git a/internal/dispatcher/config.go b/internal/dispatcher/config.go index 8f8bef2..ddc72b3 100644 --- a/internal/dispatcher/config.go +++ b/internal/dispatcher/config.go @@ -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) diff --git a/internal/dispatcher/config_test.go b/internal/dispatcher/config_test.go index 235396b..cb32f9f 100644 --- a/internal/dispatcher/config_test.go +++ b/internal/dispatcher/config_test.go @@ -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) diff --git a/internal/dispatcher/control.go b/internal/dispatcher/control.go index b5bb107..73e71da 100644 --- a/internal/dispatcher/control.go +++ b/internal/dispatcher/control.go @@ -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 } diff --git a/internal/dispatcher/control_test.go b/internal/dispatcher/control_test.go index 25dac66..12c57fe 100644 --- a/internal/dispatcher/control_test.go +++ b/internal/dispatcher/control_test.go @@ -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" { diff --git a/internal/dispatcher/runtime.go b/internal/dispatcher/runtime.go index 89c7162..b5ae468 100644 --- a/internal/dispatcher/runtime.go +++ b/internal/dispatcher/runtime.go @@ -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) diff --git a/internal/dispatcher/runtime_integration_test.go b/internal/dispatcher/runtime_integration_test.go index 670afd3..6646c24 100644 --- a/internal/dispatcher/runtime_integration_test.go +++ b/internal/dispatcher/runtime_integration_test.go @@ -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. diff --git a/internal/dispatcher/sip_conflict_test.go b/internal/dispatcher/sip_conflict_test.go index 099d888..91e8697 100644 --- a/internal/dispatcher/sip_conflict_test.go +++ b/internal/dispatcher/sip_conflict_test.go @@ -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) diff --git a/internal/dispatcher/sip_reload.go b/internal/dispatcher/sip_reload.go index 3f9eb11..6e27ad6 100644 --- a/internal/dispatcher/sip_reload.go +++ b/internal/dispatcher/sip_reload.go @@ -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)) } diff --git a/internal/dispatcher/sip_runtime_test.go b/internal/dispatcher/sip_runtime_test.go index c53d639..b13f636 100644 --- a/internal/dispatcher/sip_runtime_test.go +++ b/internal/dispatcher/sip_runtime_test.go @@ -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) + } } diff --git a/internal/store/sip.go b/internal/store/sip.go index b21fbda..4336ced 100644 --- a/internal/store/sip.go +++ b/internal/store/sip.go @@ -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. diff --git a/internal/store/sip_change_test.go b/internal/store/sip_change_test.go index d227162..4e666de 100644 --- a/internal/store/sip_change_test.go +++ b/internal/store/sip_change_test.go @@ -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)