Reject SDK-inexpressible task snapshots before admission
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user