diff --git a/internal/dispatcher/current_ai.go b/internal/dispatcher/current_ai.go new file mode 100644 index 0000000..26b8e09 --- /dev/null +++ b/internal/dispatcher/current_ai.go @@ -0,0 +1,18 @@ +package dispatcher + +import ( + "fmt" + + "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/configread" +) + +// validateCurrentAISnapshot is the shared admission gate for every source of +// approved task configuration. SDK-inexpressible settings must fail before a +// call reserves capacity or reaches the Agent; no defaults replace them. +func validateCurrentAISnapshot(snapshot configread.CurrentSnapshot) error { + if _, err := ai.BindCurrent(snapshot.Task, snapshot.Providers); err != nil { + return fmt.Errorf("task %q approved AI snapshot: %w", snapshot.Task.TaskID, err) + } + return nil +} diff --git a/internal/dispatcher/current_ai_test.go b/internal/dispatcher/current_ai_test.go new file mode 100644 index 0000000..5665def --- /dev/null +++ b/internal/dispatcher/current_ai_test.go @@ -0,0 +1,113 @@ +package dispatcher + +import ( + "bytes" + "context" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/store" +) + +func unsupportedCurrentASRTask(t *testing.T) []byte { + t.Helper() + original := currentConfigExample(t, "config-read-task-asr") + modified := bytes.Replace(original, []byte(`"language":"zh-CN"`), []byte(`"language":"fr-FR"`), 1) + if bytes.Equal(original, modified) { + t.Fatal("ASR fixture did not contain the expected language") + } + return modified +} + +func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *testing.T) { + id := "c046b893-8628-4589-ae50-619d049248a6" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + var body []byte + switch r.URL.Path { + case "/internal/v1/dispatcher/sip": + body = currentConfigExample(t, "config-read-sip") + case "/internal/v1/dispatcher/tasks": + if r.URL.Query().Get("after") == "" { + body = currentConfigExample(t, "task-discovery-page") + } else { + body = currentConfigExample(t, "task-discovery-end") + } + case "/internal/v1/dispatcher/ai-providers": + body = currentConfigExample(t, "config-read-providers") + case "/internal/v1/dispatcher/task/task-asr": + body = unsupportedCurrentASRTask(t) + case "/internal/v1/dispatcher/tenant/1001/quota": + body = currentConfigExample(t, "config-read-quota") + default: + http.Error(w, "unknown path", http.StatusNotFound) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(body) + })) + defer server.Close() + client, err := configread.NewClient(server.URL, id, "test-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + drained := false + b := CurrentBootstrap{ + DispatcherID: id, Client: client, Store: db, + VerifySIP: func(context.Context, configread.CurrentSIP) error { return nil }, + DrainControls: func(context.Context) error { drained = true; return nil }, + } + err = b.Run(context.Background()) + if err == nil || !strings.Contains(err.Error(), "ASR language") || drained { + t.Fatalf("SDK-unsupported task must not be admitted or drained: %v, drained=%t", err, drained) + } + if admitted, err := db.CanAdmit(id, 1001, "task-asr"); err != nil || admitted { + t.Fatalf("invalid AI task opened admission: admitted=%t err=%v", admitted, err) + } + if _, err := db.ReadSnapshot(id, 1001, "task-asr"); err == nil { + t.Fatal("unsupported AI snapshot was persisted as runnable") + } +} + +func TestCurrentExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing.T) { + db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + snapshot := currentPolicySnapshot(t) + if err := snapshot.Task.UnmarshalJSON(unsupportedCurrentASRTask(t)); err != nil { + t.Fatal(err) + } + id := snapshot.Task.DispatcherID + if err := db.ApplyDiscoverySnapshot(id, []configread.CurrentDiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { + t.Fatal(err) + } + if err := db.SaveSnapshot(snapshot); err != nil { + t.Fatal(err) + } + if err := db.MarkReadyForSIP(id, snapshot.SIP.Revision); err != nil { + t.Fatal(err) + } + originator := ¤tFakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} + controller := &CurrentExecuteController{ + DispatcherID: id, Store: db, Originator: originator, Publisher: ¤tFakePublisher{}, + Now: func() time.Time { return currentMonday(9, 30) }, + } + err = controller.ProcessExecute(context.Background(), currentExecuteBody(t, "bad-ai", "15003164745")) + if err == nil || !strings.Contains(err.Error(), "ASR language") || len(originator.calls) != 0 { + t.Fatalf("corrupt AI config must not reach Agent: err=%v calls=%d", err, len(originator.calls)) + } + if admitted, err := db.CanAdmit(id, 1001, "task-asr"); err != nil || admitted { + t.Fatalf("corrupt persisted task remained eligible: admitted=%t err=%v", admitted, err) + } +} diff --git a/internal/dispatcher/current_config.go b/internal/dispatcher/current_config.go index 424dc79..93fcec6 100644 --- a/internal/dispatcher/current_config.go +++ b/internal/dispatcher/current_config.go @@ -72,6 +72,9 @@ func (b CurrentBootstrap) Run(ctx context.Context) error { 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 err := validateCurrentAISnapshot(snapshot); err != nil { + return err + } if err := b.Store.SaveSnapshot(snapshot); err != nil { return fmt.Errorf("persist immutable task %q: %w", task.TaskID, err) } diff --git a/internal/dispatcher/current_control.go b/internal/dispatcher/current_control.go index e2adf66..e652bd8 100644 --- a/internal/dispatcher/current_control.go +++ b/internal/dispatcher/current_control.go @@ -93,6 +93,9 @@ func (c *CurrentControlController) ProcessControl(ctx context.Context, body []by if err := c.VerifySIP(ctx, snapshot.SIP); err != nil { return fmt.Errorf("resume SIP revision not applied: %w", err) } + if err := validateCurrentAISnapshot(snapshot); err != nil { + return err + } if err := c.Store.SaveSnapshot(snapshot); err != nil { return fmt.Errorf("bind fresh resume task: %w", err) } diff --git a/internal/dispatcher/current_discovery.go b/internal/dispatcher/current_discovery.go index 387095d..e29c5ca 100644 --- a/internal/dispatcher/current_discovery.go +++ b/internal/dispatcher/current_discovery.go @@ -63,6 +63,9 @@ func (f *CurrentDiscoveryFollower) Poll(ctx context.Context) error { if err := f.VerifySIP(ctx, snapshot.SIP); err != nil { return f.fail(fmt.Errorf("task %q SIP is not loaded by Agent/Asterisk: %w", task.TaskID, err)) } + if err := validateCurrentAISnapshot(snapshot); err != nil { + return f.fail(err) + } if err := f.Store.SaveSnapshot(snapshot); err != nil { return f.fail(fmt.Errorf("bind discovered task %q: %w", task.TaskID, err)) } diff --git a/internal/dispatcher/current_execute.go b/internal/dispatcher/current_execute.go index 4ba4d25..db8a98e 100644 --- a/internal/dispatcher/current_execute.go +++ b/internal/dispatcher/current_execute.go @@ -139,6 +139,9 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor if err != nil { return fmt.Errorf("load authorized task %q: %w", cmd.TaskID, err) } + if err := validateCurrentAISnapshot(snapshot); err != nil { + return errors.Join(err, c.Store.CloseAdmission(c.DispatcherID)) + } loaded, err := c.Originator.LoadedTrunks(ctx) if err != nil { return fmt.Errorf("read applied Agent/Asterisk SIP revisions: %w", err) diff --git a/internal/dispatcher/current_sip_reload.go b/internal/dispatcher/current_sip_reload.go index 78ba832..6454d9c 100644 --- a/internal/dispatcher/current_sip_reload.go +++ b/internal/dispatcher/current_sip_reload.go @@ -93,6 +93,9 @@ func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscov 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 := validateCurrentAISnapshot(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)) }