From 13df3c2c2f25f30c5aa463a379eedf2879a02c50 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 12:05:35 +0800 Subject: [PATCH] Normalize Dispatcher runtime and policy code names --- .../current_dispatcher_command.go | 8 +- .../saas-dispatcher-implementation.md | 1 + internal/dispatcher/{current_ai.go => ai.go} | 4 +- .../{current_ai_test.go => ai_test.go} | 36 ++++----- internal/dispatcher/approved_control.go | 4 +- internal/dispatcher/approved_control_test.go | 18 ++--- internal/dispatcher/approved_originator.go | 6 +- .../dispatcher/approved_originator_test.go | 6 +- .../{current_config.go => config.go} | 10 +-- ...{current_config_test.go => config_test.go} | 34 ++++---- .../{current_control.go => control.go} | 18 ++--- ...urrent_control_test.go => control_test.go} | 80 +++++++++---------- .../{current_discovery.go => discovery.go} | 10 +-- ...nt_discovery_test.go => discovery_test.go} | 22 ++--- .../{current_execute.go => execute.go} | 42 +++++----- ...urrent_execute_test.go => execute_test.go} | 54 ++++++------- .../{current_gate_test.go => gate_test.go} | 12 +-- .../{current_policy.go => policy.go} | 60 +++++++------- ...{current_policy_test.go => policy_test.go} | 48 +++++------ .../{current_runtime.go => runtime.go} | 32 ++++---- ...on_test.go => runtime_integration_test.go} | 58 +++++++------- ..._conflict_test.go => sip_conflict_test.go} | 16 ++-- .../{current_sip_reload.go => sip_reload.go} | 6 +- ...ip_runtime_test.go => sip_runtime_test.go} | 22 ++--- .../rpc/approved_control_transport_test.go | 2 +- .../rpc/approved_full_ai_integration_test.go | 2 +- internal/rpc/approved_integration_test.go | 2 +- 27 files changed, 307 insertions(+), 306 deletions(-) rename internal/dispatcher/{current_ai.go => ai.go} (75%) rename internal/dispatcher/{current_ai_test.go => ai_test.go} (74%) rename internal/dispatcher/{current_config.go => config.go} (90%) rename internal/dispatcher/{current_config_test.go => config_test.go} (82%) rename internal/dispatcher/{current_control.go => control.go} (90%) rename internal/dispatcher/{current_control_test.go => control_test.go} (66%) rename internal/dispatcher/{current_discovery.go => discovery.go} (89%) rename internal/dispatcher/{current_discovery_test.go => discovery_test.go} (73%) rename internal/dispatcher/{current_execute.go => execute.go} (86%) rename internal/dispatcher/{current_execute_test.go => execute_test.go} (70%) rename internal/dispatcher/{current_gate_test.go => gate_test.go} (83%) rename internal/dispatcher/{current_policy.go => policy.go} (62%) rename internal/dispatcher/{current_policy_test.go => policy_test.go} (55%) rename internal/dispatcher/{current_runtime.go => runtime.go} (88%) rename internal/dispatcher/{current_runtime_integration_test.go => runtime_integration_test.go} (81%) rename internal/dispatcher/{current_sip_conflict_test.go => sip_conflict_test.go} (76%) rename internal/dispatcher/{current_sip_reload.go => sip_reload.go} (95%) rename internal/dispatcher/{current_sip_runtime_test.go => sip_runtime_test.go} (73%) diff --git a/cmd/sip-go-agent/current_dispatcher_command.go b/cmd/sip-go-agent/current_dispatcher_command.go index 80da402..71d1b10 100644 --- a/cmd/sip-go-agent/current_dispatcher_command.go +++ b/cmd/sip-go-agent/current_dispatcher_command.go @@ -128,15 +128,15 @@ func runCurrentDispatcher(ctx context.Context, mode string) (result error) { return coordinator.ApprovedMeta(callCtx, agentEndpoint.AgentID) }, } - worker := &dispatcher.CurrentRuntime{ + worker := &dispatcher.Runtime{ Broker: broker, - Bootstrap: dispatcher.CurrentBootstrap{ + Bootstrap: dispatcher.Bootstrap{ DispatcherID: settings.DispatcherID, Client: reader, Store: database, VerifySIP: originator.VerifySIP, }, - Execute: dispatcher.CurrentExecuteController{ + Execute: dispatcher.ExecuteController{ DispatcherID: settings.DispatcherID, Store: database, Originator: originator, Publisher: broker, Now: time.Now, }, - Control: dispatcher.CurrentControlController{ + Control: dispatcher.ControlController{ DispatcherID: settings.DispatcherID, Store: database, Client: reader, Agent: originator, VerifySIP: originator.VerifySIP, Now: time.Now, }, PollInterval: time.Second, DiscoveryInterval: time.Minute, Logger: slog.Default(), diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 89b93ef..c309486 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -109,6 +109,7 @@ - MQ/租户路由名称收敛:现行 MQ 和租户路由的自有 `Current*` 代码标识改为唯一的 `Broker`/`Open`、`Route` 及控制/任务/结果路由入口;七份 Go 文件改为无代次路径。固定 Topic/队列的 `.v1` 名称和值保持原样,隔离 MQ 的精确路由和不声明拓扑检查继续执行。 - HTTP 配置读取旧入口:移除旧字符串租户的三资源拼接读取、单独任务状态读取及旧 snapshot/changes 发现分支和专属测试;保留统一的凭据头、HTTP/JSON 错误处理和当前五类只读读取。当前任务/发现的 503、身份错配及分页测试继续执行,不从旧配置回退。 - HTTP 配置读取名称收敛:四份现行读取/发现 Go 文件改为 `snapshots`/`discovery` 的无代次路径,公开快照与 `ReadSIP`/`ReadTask`/`ReadTasks`/`ReadAllTasks` 只保留一套入口。原始 AI JSON、租户数字身份、五类 HTTP 路径及断连拒绝行为未改变;相关调用方和隔离测试同步更新。 +- Dispatcher 名称收敛:现行 `Current*` 调度、控制、发现、线路选择及执行类型/错误改为 `Runtime`、`Bootstrap`、`DiscoveryFollower` 等唯一代码入口;18 份 Go 源码/测试文件移除代次路径,相关 Agent 控制及调用方同步更新。任务归属、固定 SIP 快照、白名单/时段、额度和结果防重规则不变,现行隔离执行测试继续通过。 ## 验收台账 diff --git a/internal/dispatcher/current_ai.go b/internal/dispatcher/ai.go similarity index 75% rename from internal/dispatcher/current_ai.go rename to internal/dispatcher/ai.go index 3ad747c..98d863a 100644 --- a/internal/dispatcher/current_ai.go +++ b/internal/dispatcher/ai.go @@ -7,10 +7,10 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -// validateCurrentAISnapshot is the shared admission gate for every source of +// validateAISnapshot 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.Snapshot) error { +func validateAISnapshot(snapshot configread.Snapshot) error { if _, err := ai.BindCurrent(snapshot.Task, snapshot.Providers); err != nil { return fmt.Errorf("task %q approved AI snapshot: %w", snapshot.Task.TaskID, err) } diff --git a/internal/dispatcher/current_ai_test.go b/internal/dispatcher/ai_test.go similarity index 74% rename from internal/dispatcher/current_ai_test.go rename to internal/dispatcher/ai_test.go index 3759f9e..9ff3e2b 100644 --- a/internal/dispatcher/current_ai_test.go +++ b/internal/dispatcher/ai_test.go @@ -14,9 +14,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -func unsupportedCurrentASRTask(t *testing.T) []byte { +func unsupportedASRTask(t *testing.T) []byte { t.Helper() - original := currentConfigExample(t, "config-read-task-asr") + original := configExample(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") @@ -24,25 +24,25 @@ func unsupportedCurrentASRTask(t *testing.T) []byte { return modified } -func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *testing.T) { +func TestBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(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") + body = configExample(t, "config-read-sip") case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - body = currentConfigExample(t, "task-discovery-page") + body = configExample(t, "task-discovery-page") } else { - body = currentConfigExample(t, "task-discovery-end") + body = configExample(t, "task-discovery-end") } case "/internal/v1/dispatcher/ai-providers": - body = currentConfigExample(t, "config-read-providers") + body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - body = unsupportedCurrentASRTask(t) + body = unsupportedASRTask(t) case "/internal/v1/dispatcher/tenant/1001/quota": - body = currentConfigExample(t, "config-read-quota") + body = configExample(t, "config-read-quota") default: http.Error(w, "unknown path", http.StatusNotFound) return @@ -61,7 +61,7 @@ func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *test } defer db.Close() drained := false - b := CurrentBootstrap{ + b := Bootstrap{ DispatcherID: id, Client: client, Store: db, VerifySIP: func(context.Context, configread.SIP) error { return nil }, DrainControls: func(context.Context) error { drained = true; return nil }, @@ -78,14 +78,14 @@ func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *test } } -func TestCurrentExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing.T) { +func TestExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing.T) { db, err := store.Open(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 { + snapshot := policySnapshot(t) + if err := snapshot.Task.UnmarshalJSON(unsupportedASRTask(t)); err != nil { t.Fatal(err) } id := snapshot.Task.DispatcherID @@ -98,12 +98,12 @@ func TestCurrentExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing. 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) }, + originator := &fakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} + controller := &ExecuteController{ + DispatcherID: id, Store: db, Originator: originator, Publisher: &fakePublisher{}, + Now: func() time.Time { return monday(9, 30) }, } - err = controller.ProcessExecute(context.Background(), currentExecuteBody(t, "bad-ai", "15003164745")) + err = controller.ProcessExecute(context.Background(), executeBody(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)) } diff --git a/internal/dispatcher/approved_control.go b/internal/dispatcher/approved_control.go index 6ebab43..ddf3afd 100644 --- a/internal/dispatcher/approved_control.go +++ b/internal/dispatcher/approved_control.go @@ -10,12 +10,12 @@ import ( "github.com/google/uuid" ) -var _ CurrentControlAgent = (*ApprovedOriginator)(nil) +var _ ControlAgent = (*ApprovedOriginator)(nil) // SendControl dispatches a task-level instruction using the same authenticated // Agent session as approved execution. Each delivery gets an internal RPC trace // ID, but no external command ID, idempotency key or revision CAS is invented. -func (o *ApprovedOriginator) SendControl(ctx context.Context, spec CurrentControlSpec) error { +func (o *ApprovedOriginator) SendControl(ctx context.Context, spec ControlSpec) error { if ctx == nil || o == nil || spec.DispatcherID == "" || spec.DispatcherID != o.DispatcherID || spec.TenantID <= 0 || strings.TrimSpace(spec.TaskID) == "" { return errors.New("approved Agent task control identity is incomplete") } diff --git a/internal/dispatcher/approved_control_test.go b/internal/dispatcher/approved_control_test.go index c06ee81..e5ddb9e 100644 --- a/internal/dispatcher/approved_control_test.go +++ b/internal/dispatcher/approved_control_test.go @@ -24,7 +24,7 @@ func TestApprovedAgentControlUsesTaskIdentityAndNoExternalCommandDedup(t *testin {"stop", "drain", agentpb.ControlAction_CONTROL_ACTION_STOP, agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN}, {"resume", "", agentpb.ControlAction_CONTROL_ACTION_RESUME, agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_UNSPECIFIED}, } { - spec := CurrentControlSpec{DispatcherID: originator.DispatcherID, TenantID: execution.TenantID, TaskID: execution.TaskID, Action: tc.action, ActiveCallPolicy: tc.policy, Reason: "synthetic-private-reason-not-for-agent"} + spec := ControlSpec{DispatcherID: originator.DispatcherID, TenantID: execution.TenantID, TaskID: execution.TaskID, Action: tc.action, ActiveCallPolicy: tc.policy, Reason: "synthetic-private-reason-not-for-agent"} if err := originator.SendControl(context.Background(), spec); err != nil { t.Fatalf("current %s control was not applied: %v", tc.action, err) } @@ -48,17 +48,17 @@ func TestApprovedAgentControlUsesTaskIdentityAndNoExternalCommandDedup(t *testin func TestApprovedAgentControlRejectsInvalidInstructionsAndUnknownReply(t *testing.T) { fake := &fakeApprovedAgentRPC{} originator, execution := approvedOriginatorFixture(t, fake) - base := CurrentControlSpec{DispatcherID: originator.DispatcherID, TenantID: execution.TenantID, TaskID: execution.TaskID, Action: "pause", ActiveCallPolicy: "hangup"} + base := ControlSpec{DispatcherID: originator.DispatcherID, TenantID: execution.TenantID, TaskID: execution.TaskID, Action: "pause", ActiveCallPolicy: "hangup"} for _, tc := range []struct { name string - edit func(*CurrentControlSpec) + edit func(*ControlSpec) }{ - {"wrong dispatcher", func(spec *CurrentControlSpec) { spec.DispatcherID = "another-dispatcher" }}, - {"tenant zero", func(spec *CurrentControlSpec) { spec.TenantID = 0 }}, - {"task missing", func(spec *CurrentControlSpec) { spec.TaskID = "" }}, - {"unsupported action", func(spec *CurrentControlSpec) { spec.Action = "restart" }}, - {"pause missing policy", func(spec *CurrentControlSpec) { spec.ActiveCallPolicy = "" }}, - {"resume with hangup", func(spec *CurrentControlSpec) { spec.Action = "resume" }}, + {"wrong dispatcher", func(spec *ControlSpec) { spec.DispatcherID = "another-dispatcher" }}, + {"tenant zero", func(spec *ControlSpec) { spec.TenantID = 0 }}, + {"task missing", func(spec *ControlSpec) { spec.TaskID = "" }}, + {"unsupported action", func(spec *ControlSpec) { spec.Action = "restart" }}, + {"pause missing policy", func(spec *ControlSpec) { spec.ActiveCallPolicy = "" }}, + {"resume with hangup", func(spec *ControlSpec) { spec.Action = "resume" }}, } { t.Run(tc.name, func(t *testing.T) { spec := base diff --git a/internal/dispatcher/approved_originator.go b/internal/dispatcher/approved_originator.go index 9865388..0f60d41 100644 --- a/internal/dispatcher/approved_originator.go +++ b/internal/dispatcher/approved_originator.go @@ -30,7 +30,7 @@ type ApprovedOriginator struct { Meta func(context.Context) (*agentpb.RequestMeta, error) } -var _ CurrentOriginator = (*ApprovedOriginator)(nil) +var _ Originator = (*ApprovedOriginator)(nil) func (o *ApprovedOriginator) activeMeta(ctx context.Context) (*agentpb.RequestMeta, error) { if o == nil || o.DispatcherID == "" || o.Client == nil || o.Meta == nil { @@ -95,14 +95,14 @@ func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.SIP) return nil } -func (o *ApprovedOriginator) Originate(ctx context.Context, spec CurrentCallSpec) error { +func (o *ApprovedOriginator) Originate(ctx context.Context, spec CallSpec) error { if o == nil || spec.EventID == "" || spec.TaskID == "" || spec.TenantID <= 0 || spec.TrunkID == "" || spec.Callee == "" || spec.CallerID == "" || spec.DialedCallee == "" || spec.RingTimeoutMS <= 0 || spec.MaxCallDurationMS <= 0 || spec.Deadline.IsZero() { return errors.New("approved call decision is incomplete") } if spec.Snapshot.Task.DispatcherID != o.DispatcherID || spec.Snapshot.Task.TaskID != spec.TaskID || spec.Snapshot.Task.TenantID != spec.TenantID || spec.Snapshot.SIP.DispatcherID != o.DispatcherID || spec.Snapshot.SIP.Revision <= 0 { return errors.New("approved call snapshot identity does not match decision") } - if err := validateCurrentAISnapshot(spec.Snapshot); err != nil { + if err := validateAISnapshot(spec.Snapshot); err != nil { return err } meta, err := o.activeMeta(ctx) diff --git a/internal/dispatcher/approved_originator_test.go b/internal/dispatcher/approved_originator_test.go index f3a02af..25eb5e3 100644 --- a/internal/dispatcher/approved_originator_test.go +++ b/internal/dispatcher/approved_originator_test.go @@ -52,9 +52,9 @@ func (f *fakeApprovedAgentRPC) GetLoadedSIP(_ context.Context, _ *agentpb.GetLoa return &agentpb.GetLoadedSIPResponse{TrunkRevision: f.loaded}, nil } -func approvedOriginatorFixture(t *testing.T, client *fakeApprovedAgentRPC) (*ApprovedOriginator, CurrentCallSpec) { +func approvedOriginatorFixture(t *testing.T, client *fakeApprovedAgentRPC) (*ApprovedOriginator, CallSpec) { t.Helper() - snapshot := currentPolicySnapshot(t) + snapshot := policySnapshot(t) id := snapshot.Task.DispatcherID orig := &ApprovedOriginator{ DispatcherID: id, Client: client, @@ -62,7 +62,7 @@ func approvedOriginatorFixture(t *testing.T, client *fakeApprovedAgentRPC) (*App return &agentpb.RequestMeta{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1, OperationId: "read-sip"}, nil }, } - spec := CurrentCallSpec{ + spec := CallSpec{ EventID: "event-1", TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, Callee: "15003164745", DialedCallee: "708915003164745", CallerID: "BD93205882", TrunkID: "trunk-mock", RingTimeoutMS: snapshot.Task.RingTimeoutMS, diff --git a/internal/dispatcher/current_config.go b/internal/dispatcher/config.go similarity index 90% rename from internal/dispatcher/current_config.go rename to internal/dispatcher/config.go index 0ca5868..8f8bef2 100644 --- a/internal/dispatcher/current_config.go +++ b/internal/dispatcher/config.go @@ -11,20 +11,20 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -// CurrentBootstrap is the one admission boundary for the current read-only +// Bootstrap is the one admission boundary for the current read-only // configuration contract. Control consumption remains available even when // opening task admission fails. -type CurrentBootstrap struct { +type Bootstrap struct { Client *configread.Client Store *store.Store DispatcherID string VerifySIP func(context.Context, configread.SIP) error DrainControls func(context.Context) error - Cursor *string // memory-only; restart always starts from a full snapshot + Cursor *string // memory-only; restart always starts from a full snapshot SIP *configread.SIP // the exact approved snapshot verified on startup } -func (b CurrentBootstrap) Run(ctx context.Context) error { +func (b Bootstrap) Run(ctx context.Context) error { if b.Client == nil || b.Store == nil || b.DispatcherID == "" || b.VerifySIP == nil || b.DrainControls == nil { return errors.New("current Dispatcher bootstrap requires config client, durable store, identity, applied SIP verifier, and control drain") } @@ -72,7 +72,7 @@ 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 { + if err := validateAISnapshot(snapshot); err != nil { return err } if err := b.Store.SaveSnapshot(snapshot); err != nil { diff --git a/internal/dispatcher/current_config_test.go b/internal/dispatcher/config_test.go similarity index 82% rename from internal/dispatcher/current_config_test.go rename to internal/dispatcher/config_test.go index a80e792..235396b 100644 --- a/internal/dispatcher/current_config_test.go +++ b/internal/dispatcher/config_test.go @@ -13,7 +13,7 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -func currentConfigExample(t *testing.T, name string) []byte { +func configExample(t *testing.T, name string) []byte { t.Helper() raw, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", name+".json")) if err != nil { @@ -22,7 +22,7 @@ func currentConfigExample(t *testing.T, name string) []byte { return raw } -func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { +func TestBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" called := make([]string, 0) server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -30,19 +30,19 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { var response []byte switch r.URL.Path { case "/internal/v1/dispatcher/sip": - response = currentConfigExample(t, "config-read-sip") + response = configExample(t, "config-read-sip") case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - response = currentConfigExample(t, "task-discovery-page") + response = configExample(t, "task-discovery-page") } else { - response = currentConfigExample(t, "task-discovery-end") + response = configExample(t, "task-discovery-end") } case "/internal/v1/dispatcher/ai-providers": - response = currentConfigExample(t, "config-read-providers") + response = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - response = currentConfigExample(t, "config-read-task-asr") + response = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - response = currentConfigExample(t, "config-read-quota") + response = configExample(t, "config-read-quota") default: t.Errorf("unapproved config path %s", r.URL.Path) w.WriteHeader(http.StatusNotFound) @@ -64,7 +64,7 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { verifierCalled, drained := false, false var cursor string var approvedSIP configread.SIP - bootstrap := CurrentBootstrap{ + bootstrap := Bootstrap{ Client: client, Store: db, DispatcherID: id, Cursor: &cursor, SIP: &approvedSIP, VerifySIP: func(_ context.Context, sip configread.SIP) error { verifierCalled = true @@ -102,7 +102,7 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { } } -func TestCurrentBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { +func TestBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" for _, failure := range []string{"verify", "drain"} { t.Run(failure, func(t *testing.T) { @@ -110,19 +110,19 @@ func TestCurrentBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { var response []byte switch r.URL.Path { case "/internal/v1/dispatcher/sip": - response = currentConfigExample(t, "config-read-sip") + response = configExample(t, "config-read-sip") case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - response = currentConfigExample(t, "task-discovery-page") + response = configExample(t, "task-discovery-page") } else { - response = currentConfigExample(t, "task-discovery-end") + response = configExample(t, "task-discovery-end") } case "/internal/v1/dispatcher/ai-providers": - response = currentConfigExample(t, "config-read-providers") + response = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - response = currentConfigExample(t, "config-read-task-asr") + response = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - response = currentConfigExample(t, "config-read-quota") + response = configExample(t, "config-read-quota") default: t.Fatalf("unexpected path %s", r.URL.Path) } @@ -139,7 +139,7 @@ func TestCurrentBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { t.Fatal(err) } defer db.Close() - b := CurrentBootstrap{ + b := Bootstrap{ Client: client, Store: db, DispatcherID: id, VerifySIP: func(context.Context, configread.SIP) error { if failure == "verify" { diff --git a/internal/dispatcher/current_control.go b/internal/dispatcher/control.go similarity index 90% rename from internal/dispatcher/current_control.go rename to internal/dispatcher/control.go index c43890c..ed3a5f9 100644 --- a/internal/dispatcher/current_control.go +++ b/internal/dispatcher/control.go @@ -12,9 +12,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -// CurrentControlSpec is an Agent control instruction, not evidence that a +// ControlSpec is an Agent control instruction, not evidence that a // drain/hangup has completed. Never log a user-supplied reason or credentials. -type CurrentControlSpec struct { +type ControlSpec struct { DispatcherID string TenantID int64 TaskID string @@ -23,15 +23,15 @@ type CurrentControlSpec struct { Reason string } -type CurrentControlAgent interface { - SendControl(context.Context, CurrentControlSpec) error +type ControlAgent interface { + SendControl(context.Context, ControlSpec) error } -type CurrentControlController struct { +type ControlController struct { DispatcherID string Store *store.Store Client *configread.Client - Agent CurrentControlAgent + Agent ControlAgent VerifySIP func(context.Context, configread.SIP) error Now func() time.Time } @@ -40,7 +40,7 @@ type CurrentControlController struct { // 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 *CurrentControlController) ProcessControl(ctx context.Context, body []byte) error { +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") } @@ -93,7 +93,7 @@ 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 { + if err := validateAISnapshot(snapshot); err != nil { return err } if err := c.Store.SaveSnapshot(snapshot); err != nil { @@ -116,7 +116,7 @@ func (c *CurrentControlController) ProcessControl(ctx context.Context, body []by } return fmt.Errorf("prepare task control %q: %w", event.EventID, err) } - spec := CurrentControlSpec{ + spec := ControlSpec{ DispatcherID: event.DispatcherID, TenantID: event.TenantID, TaskID: event.Payload.TaskID, Action: event.Payload.Action, ActiveCallPolicy: policy, Reason: event.Payload.Reason, diff --git a/internal/dispatcher/current_control_test.go b/internal/dispatcher/control_test.go similarity index 66% rename from internal/dispatcher/current_control_test.go rename to internal/dispatcher/control_test.go index 10844a2..4c93b57 100644 --- a/internal/dispatcher/current_control_test.go +++ b/internal/dispatcher/control_test.go @@ -17,13 +17,13 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -type currentFakeControlAgent struct { - calls []CurrentControlSpec +type fakeControlAgent struct { + calls []ControlSpec err error - before func(CurrentControlSpec) + before func(ControlSpec) } -func (a *currentFakeControlAgent) SendControl(_ context.Context, spec CurrentControlSpec) error { +func (a *fakeControlAgent) SendControl(_ context.Context, spec ControlSpec) error { if a.before != nil { a.before(spec) } @@ -31,10 +31,10 @@ func (a *currentFakeControlAgent) SendControl(_ context.Context, spec CurrentCon return a.err } -func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *currentFakeControlAgent, *store.Store) { +func newControlFixture(t *testing.T) (*ControlController, *fakeControlAgent, *store.Store) { t.Helper() id := "c046b893-8628-4589-ae50-619d049248a6" - snapshot := currentPolicySnapshot(t) + snapshot := policySnapshot(t) approvedSIP, err := json.Marshal(snapshot.SIP) if err != nil { t.Fatal(err) @@ -46,11 +46,11 @@ func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *current case "/internal/v1/dispatcher/sip": body = approvedSIP case "/internal/v1/dispatcher/ai-providers": - body = currentConfigExample(t, "config-read-providers") + body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - body = currentConfigExample(t, "config-read-task-asr") + body = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - body = currentConfigExample(t, "config-read-quota") + body = configExample(t, "config-read-quota") default: t.Errorf("unexpected HTTP path %s", r.URL.Path) w.WriteHeader(http.StatusNotFound) @@ -77,14 +77,14 @@ func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *current if err := s.MarkReadyForSIP(id, snapshot.SIP.Revision); err != nil { t.Fatal(err) } - agent := ¤tFakeControlAgent{} - controller := &CurrentControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, VerifySIP: func(context.Context, configread.SIP) error { return nil }, Now: func() time.Time { return currentMonday(9, 30) }} + 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) }} return controller, agent, s } -func currentControlBody(t *testing.T, eventID, action, policy string) []byte { +func controlBody(t *testing.T, eventID, action, policy string) []byte { t.Helper() - body := string(currentConfigExample(t, "mq-control")) + body := string(configExample(t, "mq-control")) body = strings.Replace(body, `"event_id":"control-example"`, `"event_id":"`+eventID+`"`, 1) body = strings.Replace(body, `"action":"pause"`, `"action":"`+action+`"`, 1) body = strings.Replace(body, `"active_call_policy":"drain"`, `"active_call_policy":"`+policy+`"`, 1) @@ -94,9 +94,9 @@ func currentControlBody(t *testing.T, eventID, action, policy string) []byte { return []byte(body) } -func TestCurrentControlPauseDefaultHangupAckAfterAgentDispatch(t *testing.T) { - controller, agent, s := newCurrentControlFixture(t) - agent.before = func(spec CurrentControlSpec) { +func TestControlPauseDefaultHangupAckAfterAgentDispatch(t *testing.T) { + controller, agent, s := newControlFixture(t) + agent.before = func(spec ControlSpec) { if spec.Action != "pause" || spec.ActiveCallPolicy != "hangup" { t.Errorf("wrong default control policy: %+v", spec) } @@ -107,7 +107,7 @@ func TestCurrentControlPauseDefaultHangupAckAfterAgentDispatch(t *testing.T) { t.Errorf("ack written before Agent dispatch: %+v %v", outbox, err) } } - body := currentControlBody(t, "pause-1", "pause", "") + body := controlBody(t, "pause-1", "pause", "") if err := controller.ProcessControl(context.Background(), body); err != nil { t.Fatal(err) } @@ -131,13 +131,13 @@ func TestCurrentControlPauseDefaultHangupAckAfterAgentDispatch(t *testing.T) { } } -func TestCurrentControlStopCannotResumeOrDispatchPending(t *testing.T) { - controller, agent, s := newCurrentControlFixture(t) +func TestControlStopCannotResumeOrDispatchPending(t *testing.T) { + controller, agent, s := newControlFixture(t) command := store.ExecuteCommand{DispatcherID: controller.DispatcherID, EventID: "queued-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:00:00Z"} if _, _, err := s.RecordExecute(command); err != nil { t.Fatal(err) } - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "stop-1", "stop", "")); err != nil { + if err := controller.ProcessControl(context.Background(), controlBody(t, "stop-1", "stop", "")); err != nil { t.Fatal(err) } if len(agent.calls) != 1 || agent.calls[0].ActiveCallPolicy != "hangup" { @@ -146,7 +146,7 @@ func TestCurrentControlStopCannotResumeOrDispatchPending(t *testing.T) { if pending, err := s.ListPendingExecute(controller.DispatcherID); err != nil || len(pending) != 0 { t.Fatalf("old pending call survived stop: %+v %v", pending, err) } - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "resume-1", "resume", "")); err != nil { + if err := controller.ProcessControl(context.Background(), controlBody(t, "resume-1", "resume", "")); err != nil { t.Fatal(err) } if len(agent.calls) != 1 { @@ -158,14 +158,14 @@ func TestCurrentControlStopCannotResumeOrDispatchPending(t *testing.T) { } } -func TestCurrentStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox(t *testing.T) { - control, _, s := newCurrentControlFixture(t) - if err := control.ProcessControl(context.Background(), currentControlBody(t, "stop-before-execute", "stop", "")); err != nil { +func TestStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox(t *testing.T) { + control, _, s := newControlFixture(t) + if err := control.ProcessControl(context.Background(), controlBody(t, "stop-before-execute", "stop", "")); err != nil { t.Fatal(err) } - originator := ¤tFakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} - execute := &CurrentExecuteController{DispatcherID: control.DispatcherID, Store: s, Originator: originator, Publisher: ¤tFakePublisher{}, Now: func() time.Time { return currentMonday(9, 30) }} - if err := execute.ProcessExecute(context.Background(), currentExecuteBody(t, "old-queued-1", "15003164745")); err != nil { + originator := &fakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} + execute := &ExecuteController{DispatcherID: control.DispatcherID, Store: s, Originator: originator, Publisher: &fakePublisher{}, Now: func() time.Time { return monday(9, 30) }} + if err := execute.ProcessExecute(context.Background(), executeBody(t, "old-queued-1", "15003164745")); err != nil { t.Fatal(err) } if len(originator.calls) != 0 { @@ -179,16 +179,16 @@ func TestCurrentStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox( } } -func TestCurrentResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) { - controller, agent, s := newCurrentControlFixture(t) - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "pause-before-resume", "pause", "drain")); err != nil { +func TestResumeRequiresFreshTaskAndAppliedSIP(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 := controller.ProcessControl(context.Background(), currentControlBody(t, "resume-after-load", "resume", "")); err == nil { + 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 { @@ -198,7 +198,7 @@ func TestCurrentResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) { t.Fatalf("failed SIP check reopened admission: %v %v", admitted, err) } controller.VerifySIP = approved - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "resume-after-load", "resume", "")); err != nil { + 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 != "" { @@ -213,11 +213,11 @@ func TestCurrentResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) { } } -func TestCurrentRuntimeSurfacesControlFailureToCloseAdmission(t *testing.T) { - controller, agent, s := newCurrentControlFixture(t) +func TestRuntimeSurfacesControlFailureToCloseAdmission(t *testing.T) { + controller, agent, s := newControlFixture(t) agent.err = errors.New("injected Agent control failure") - runtime := &CurrentRuntime{Bootstrap: CurrentBootstrap{DispatcherID: controller.DispatcherID, Store: s}, Control: *controller, Logger: slog.Default(), failures: make(chan error, 1), taskLocks: make(map[string]*sync.Mutex)} - if err := runtime.handleControl(context.Background(), "", currentControlBody(t, "fatal-control", "pause", "")); err == nil { + runtime := &Runtime{Bootstrap: Bootstrap{DispatcherID: controller.DispatcherID, Store: s}, Control: *controller, Logger: slog.Default(), failures: make(chan error, 1), taskLocks: make(map[string]*sync.Mutex)} + if err := runtime.handleControl(context.Background(), "", controlBody(t, "fatal-control", "pause", "")); err == nil { t.Fatal("Agent control failure was hidden") } select { @@ -230,10 +230,10 @@ func TestCurrentRuntimeSurfacesControlFailureToCloseAdmission(t *testing.T) { } } -func TestCurrentControlAgentFailureKeepsBarrierWithoutSuccessAck(t *testing.T) { - controller, agent, s := newCurrentControlFixture(t) +func TestControlAgentFailureKeepsBarrierWithoutSuccessAck(t *testing.T) { + controller, agent, s := newControlFixture(t) agent.err = errors.New("injected Agent control RPC timeout") - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "stop-timeout", "stop", "drain")); err == nil { + if err := controller.ProcessControl(context.Background(), controlBody(t, "stop-timeout", "stop", "drain")); err == nil { t.Fatal("Agent control failure was hidden") } if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted { @@ -243,7 +243,7 @@ func TestCurrentControlAgentFailureKeepsBarrierWithoutSuccessAck(t *testing.T) { t.Fatalf("false applied ack after Agent failure: %+v %v", outbox, err) } agent.err = nil - if err := controller.ProcessControl(context.Background(), currentControlBody(t, "stop-timeout", "stop", "drain")); err != nil || len(agent.calls) != 2 { + if err := controller.ProcessControl(context.Background(), controlBody(t, "stop-timeout", "stop", "drain")); err != nil || len(agent.calls) != 2 { t.Fatalf("redelivered control not dispatched: %+v %v", agent.calls, err) } } diff --git a/internal/dispatcher/current_discovery.go b/internal/dispatcher/discovery.go similarity index 89% rename from internal/dispatcher/current_discovery.go rename to internal/dispatcher/discovery.go index bab3290..e0142ea 100644 --- a/internal/dispatcher/current_discovery.go +++ b/internal/dispatcher/discovery.go @@ -11,9 +11,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -// CurrentDiscoveryFollower keeps only an in-memory cursor. A restart always +// DiscoveryFollower keeps only an in-memory cursor. A restart always // takes a complete snapshot and drains the control queue before admitting. -type CurrentDiscoveryFollower struct { +type DiscoveryFollower struct { DispatcherID string Client *configread.Client Store *store.Store @@ -22,7 +22,7 @@ type CurrentDiscoveryFollower struct { Cursor string } -func (f *CurrentDiscoveryFollower) Poll(ctx context.Context) error { +func (f *DiscoveryFollower) Poll(ctx context.Context) error { if f == nil || f.Client == nil || f.Store == nil || f.DispatcherID == "" || f.ApprovedSIP.DispatcherID != f.DispatcherID || f.ApprovedSIP.Revision <= 0 || f.VerifySIP == nil || f.Cursor == "" { return errors.New("discovery follower requires a verified full snapshot and in-memory cursor") } @@ -63,7 +63,7 @@ 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 { + if err := validateAISnapshot(snapshot); err != nil { return f.fail(err) } if err := f.Store.SaveSnapshot(snapshot); err != nil { @@ -74,7 +74,7 @@ func (f *CurrentDiscoveryFollower) Poll(ctx context.Context) error { } } -func (f *CurrentDiscoveryFollower) fail(cause error) error { +func (f *DiscoveryFollower) fail(cause error) error { if err := f.Store.CloseAdmission(f.DispatcherID); err != nil { return errors.Join(cause, fmt.Errorf("close admission after discovery failure: %w", err)) } diff --git a/internal/dispatcher/current_discovery_test.go b/internal/dispatcher/discovery_test.go similarity index 73% rename from internal/dispatcher/current_discovery_test.go rename to internal/dispatcher/discovery_test.go index c3168ab..542a4ce 100644 --- a/internal/dispatcher/current_discovery_test.go +++ b/internal/dispatcher/discovery_test.go @@ -11,10 +11,10 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func newCurrentFollowerFixture(t *testing.T, revision int64) (*CurrentDiscoveryFollower, func() bool) { +func newFollowerFixture(t *testing.T, revision int64) (*DiscoveryFollower, func() bool) { t.Helper() - executor, _, _, s := newCurrentExecuteFixture(t) - approved := currentPolicySnapshot(t).SIP + executor, _, _, s := newExecuteFixture(t) + approved := policySnapshot(t).SIP sipBody, err := json.Marshal(approved) if err != nil { t.Fatal(err) @@ -39,11 +39,11 @@ func newCurrentFollowerFixture(t *testing.T, revision int64) (*CurrentDiscoveryF case "/internal/v1/dispatcher/sip": body = sipBody case "/internal/v1/dispatcher/ai-providers": - body = currentConfigExample(t, "config-read-providers") + body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - body = currentConfigExample(t, "config-read-task-asr") + body = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - body = currentConfigExample(t, "config-read-quota") + body = configExample(t, "config-read-quota") default: t.Errorf("unexpected config path: %s", r.URL.Path) w.WriteHeader(http.StatusNotFound) @@ -56,12 +56,12 @@ func newCurrentFollowerFixture(t *testing.T, revision int64) (*CurrentDiscoveryF if err != nil { t.Fatal(err) } - follower := &CurrentDiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: func(context.Context, configread.SIP) error { return nil }, Cursor: "initial"} + follower := &DiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: func(context.Context, configread.SIP) error { return nil }, Cursor: "initial"} return follower, func() bool { return terminalSeen } } -func TestCurrentDiscoveryFollowerPersistsPageBeforeAdvancingToTerminalEmpty(t *testing.T) { - follower, terminal := newCurrentFollowerFixture(t, 1) +func TestDiscoveryFollowerPersistsPageBeforeAdvancingToTerminalEmpty(t *testing.T) { + follower, terminal := newFollowerFixture(t, 1) if err := follower.Poll(context.Background()); err != nil { t.Fatal(err) } @@ -73,8 +73,8 @@ func TestCurrentDiscoveryFollowerPersistsPageBeforeAdvancingToTerminalEmpty(t *t } } -func TestCurrentDiscoveryFollowerFailedPageKeepsCursorAndClosesAdmission(t *testing.T) { - follower, terminal := newCurrentFollowerFixture(t, 2) +func TestDiscoveryFollowerFailedPageKeepsCursorAndClosesAdmission(t *testing.T) { + follower, terminal := newFollowerFixture(t, 2) if err := follower.Poll(context.Background()); err == nil { t.Fatal("task revision mismatch was accepted") } diff --git a/internal/dispatcher/current_execute.go b/internal/dispatcher/execute.go similarity index 86% rename from internal/dispatcher/current_execute.go rename to internal/dispatcher/execute.go index d364155..f1f8f0a 100644 --- a/internal/dispatcher/current_execute.go +++ b/internal/dispatcher/execute.go @@ -12,9 +12,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -// CurrentCallSpec is a frozen instruction handed to one Agent. Its Snapshot +// CallSpec is a frozen instruction handed to one Agent. Its Snapshot // includes immutable AI parameters and provider credentials; never log it. -type CurrentCallSpec struct { +type CallSpec struct { DispatcherID string EventID string TenantID int64 @@ -30,24 +30,24 @@ type CurrentCallSpec struct { Snapshot configread.Snapshot } -type CurrentOriginator interface { +type Originator interface { LoadedTrunks(context.Context) (map[string]int64, error) - Originate(context.Context, CurrentCallSpec) error + Originate(context.Context, CallSpec) error } -type CurrentPublisher interface { +type Publisher interface { Publish(context.Context, string, string, []byte) error } -type CurrentExecuteController struct { +type ExecuteController struct { DispatcherID string Store *store.Store - Originator CurrentOriginator - Publisher CurrentPublisher + Originator Originator + Publisher Publisher Now func() time.Time } -func (c *CurrentExecuteController) validate() error { +func (c *ExecuteController) validate() error { if c == nil || c.DispatcherID == "" || c.Store == nil || c.Originator == nil || c.Publisher == nil || c.Now == nil { return errors.New("current execution requires Dispatcher identity, durable store, originator, publisher, and clock") } @@ -58,7 +58,7 @@ func (c *CurrentExecuteController) validate() error { // Its business acknowledgment is a separate durable outbox event, written // only after originate was actually dispatched or an individual number was // definitively rejected. -func (c *CurrentExecuteController) ProcessExecute(ctx context.Context, body []byte) error { +func (c *ExecuteController) ProcessExecute(ctx context.Context, body []byte) error { if err := c.validate(); err != nil { return err } @@ -99,8 +99,8 @@ func (c *CurrentExecuteController) ProcessExecute(ctx context.Context, body []by if stored.Status != "pending" { return nil } - if _, whitelisted := currentCalleeWhitelist[cmd.Callee]; !whitelisted { - if err := c.Store.RejectExecute(cmd.DispatcherID, cmd.EventID, ErrCurrentCalleeRejected.Error()); err != nil { + if _, whitelisted := calleeWhitelist[cmd.Callee]; !whitelisted { + if err := c.Store.RejectExecute(cmd.DispatcherID, cmd.EventID, ErrCalleeRejected.Error()); err != nil { return fmt.Errorf("persist rejected call %q: %w", cmd.EventID, err) } return nil @@ -110,7 +110,7 @@ func (c *CurrentExecuteController) ProcessExecute(ctx context.Context, body []by // ProcessPending reevaluates only durable, never-originated commands. This is // used after rule/config changes; dispatching and unknown calls are excluded. -func (c *CurrentExecuteController) ProcessPending(ctx context.Context) error { +func (c *ExecuteController) ProcessPending(ctx context.Context) error { if err := c.validate(); err != nil { return err } @@ -127,7 +127,7 @@ func (c *CurrentExecuteController) ProcessPending(ctx context.Context) error { return errors.Join(failures...) } -func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd store.ExecuteCommand) error { +func (c *ExecuteController) dispatchPending(ctx context.Context, cmd store.ExecuteCommand) error { at := c.Now() issued, err := time.Parse(time.RFC3339Nano, cmd.IssuedAt) if err != nil { @@ -140,7 +140,7 @@ 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 { + if err := validateAISnapshot(snapshot); err != nil { return errors.Join(err, c.Store.CloseAdmission(c.DispatcherID)) } loaded, err := c.Originator.LoadedTrunks(ctx) @@ -151,8 +151,8 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor if err != nil { return err } - choice, err := SelectCurrentTrunk(snapshot, cmd.Callee, at, occupied, loaded) - if errors.Is(err, ErrCurrentRuleWait) { + choice, err := SelectTrunk(snapshot, cmd.Callee, at, occupied, loaded) + if errors.Is(err, ErrRuleWait) { return nil } if err != nil { @@ -177,13 +177,13 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor cancelErr := c.Store.CancelReservationBeforeOrigin(cmd.DispatcherID, cmd.EventID) return errors.Join(fmt.Errorf("recheck applied SIP before originate: %w", err), cancelErr) } - fresh, ruleErr := SelectCurrentTrunk(snapshot, cmd.Callee, actualAt, occupied, freshLoaded) + fresh, ruleErr := SelectTrunk(snapshot, cmd.Callee, actualAt, occupied, freshLoaded) if !actualAt.Before(choice.Deadline) || ruleErr != nil || fresh.TrunkID != choice.TrunkID || fresh.CallerID != choice.CallerID { cancelErr := c.Store.CancelReservationBeforeOrigin(cmd.DispatcherID, cmd.EventID) if cancelErr != nil { return fmt.Errorf("cancel expired/changed reservation: %w", cancelErr) } - if ruleErr != nil && !errors.Is(ruleErr, ErrCurrentRuleWait) { + if ruleErr != nil && !errors.Is(ruleErr, ErrRuleWait) { return fmt.Errorf("pre-dial task %q rule: %w", cmd.TaskID, ruleErr) } return nil @@ -196,7 +196,7 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor } return cancelErr } - spec := CurrentCallSpec{ + spec := CallSpec{ DispatcherID: cmd.DispatcherID, EventID: cmd.EventID, TenantID: cmd.TenantID, TaskID: cmd.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, CallerID: choice.CallerID, Callee: cmd.Callee, DialedCallee: choice.DialedCallee, @@ -216,7 +216,7 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor // FlushOutbox uses the same event identity/body on every retry. Confirming a // bound queue is not a SaaS application receipt; confirmed rows remain stored. -func (c *CurrentExecuteController) FlushOutbox(ctx context.Context) error { +func (c *ExecuteController) FlushOutbox(ctx context.Context) error { if err := c.validate(); err != nil { return err } diff --git a/internal/dispatcher/current_execute_test.go b/internal/dispatcher/execute_test.go similarity index 70% rename from internal/dispatcher/current_execute_test.go rename to internal/dispatcher/execute_test.go index 7427893..034e35c 100644 --- a/internal/dispatcher/current_execute_test.go +++ b/internal/dispatcher/execute_test.go @@ -14,26 +14,26 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -type currentFakeOriginator struct { +type fakeOriginator struct { loaded map[string]int64 - calls []CurrentCallSpec + calls []CallSpec err error } -func (f *currentFakeOriginator) LoadedTrunks(context.Context) (map[string]int64, error) { +func (f *fakeOriginator) LoadedTrunks(context.Context) (map[string]int64, error) { return f.loaded, nil } -func (f *currentFakeOriginator) Originate(_ context.Context, spec CurrentCallSpec) error { +func (f *fakeOriginator) Originate(_ context.Context, spec CallSpec) error { f.calls = append(f.calls, spec) return f.err } -type currentFakePublisher struct { +type fakePublisher struct { bodies [][]byte err error } -func (f *currentFakePublisher) Publish(_ context.Context, exchange, key string, body []byte) error { +func (f *fakePublisher) Publish(_ context.Context, exchange, key string, body []byte) error { if exchange != "agent-call.saas.v1" || !strings.HasSuffix(key, ".out") { return errors.New("unexpected outbound MQ route") } @@ -44,14 +44,14 @@ func (f *currentFakePublisher) Publish(_ context.Context, exchange, key string, return f.err } -func newCurrentExecuteFixture(t *testing.T) (*CurrentExecuteController, *currentFakeOriginator, *currentFakePublisher, *store.Store) { +func newExecuteFixture(t *testing.T) (*ExecuteController, *fakeOriginator, *fakePublisher, *store.Store) { t.Helper() s, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } t.Cleanup(func() { _ = s.Close() }) - snapshot := currentPolicySnapshot(t) + snapshot := policySnapshot(t) if err := s.ApplyDiscoverySnapshot(snapshot.Task.DispatcherID, []configread.DiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { t.Fatal(err) } @@ -61,23 +61,23 @@ func newCurrentExecuteFixture(t *testing.T) (*CurrentExecuteController, *current if err := s.MarkReadyForSIP(snapshot.Task.DispatcherID, snapshot.SIP.Revision); err != nil { t.Fatal(err) } - originator := ¤tFakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} - publisher := ¤tFakePublisher{} - controller := &CurrentExecuteController{DispatcherID: snapshot.Task.DispatcherID, Store: s, Originator: originator, Publisher: publisher, Now: func() time.Time { return currentMonday(9, 30) }} + originator := &fakeOriginator{loaded: map[string]int64{"trunk-mock": 8}} + publisher := &fakePublisher{} + controller := &ExecuteController{DispatcherID: snapshot.Task.DispatcherID, Store: s, Originator: originator, Publisher: publisher, Now: func() time.Time { return monday(9, 30) }} return controller, originator, publisher, s } -func currentExecuteBody(t *testing.T, eventID, callee string) []byte { +func executeBody(t *testing.T, eventID, callee string) []byte { t.Helper() - body := string(currentConfigExample(t, "mq-execute")) + body := string(configExample(t, "mq-execute")) body = strings.Replace(body, `"event_id":"call-example"`, `"event_id":"`+eventID+`"`, 1) body = strings.Replace(body, `"callee":"15003164745"`, `"callee":"`+callee+`"`, 1) return []byte(body) } -func TestCurrentExecuteRejectsInvalidCalleeWithoutStoppingTask(t *testing.T) { - controller, originator, publisher, s := newCurrentExecuteFixture(t) - if err := controller.ProcessExecute(context.Background(), currentExecuteBody(t, "bad-1", "not-a-phone")); err != nil { +func TestExecuteRejectsInvalidCalleeWithoutStoppingTask(t *testing.T) { + controller, originator, publisher, s := newExecuteFixture(t) + if err := controller.ProcessExecute(context.Background(), executeBody(t, "bad-1", "not-a-phone")); err != nil { t.Fatal(err) } if len(originator.calls) != 0 { @@ -87,13 +87,13 @@ func TestCurrentExecuteRejectsInvalidCalleeWithoutStoppingTask(t *testing.T) { if err != nil || len(outbox) != 1 || !strings.Contains(string(outbox[0].Body), `"status":"rejected"`) { t.Fatalf("missing individual rejection: %+v %v", outbox, err) } - if err := controller.ProcessExecute(context.Background(), currentExecuteBody(t, "good-1", "15003164745")); err != nil { + if err := controller.ProcessExecute(context.Background(), executeBody(t, "good-1", "15003164745")); err != nil { t.Fatal(err) } if len(originator.calls) != 1 { t.Fatalf("invalid number paused entire task; originate calls=%d", len(originator.calls)) } - if err := controller.ProcessExecute(context.Background(), currentExecuteBody(t, "good-1", "15003164745")); err != nil || len(originator.calls) != 1 { + if err := controller.ProcessExecute(context.Background(), executeBody(t, "good-1", "15003164745")); err != nil || len(originator.calls) != 1 { t.Fatalf("duplicate command reoriginated: %d %v", len(originator.calls), err) } if err := controller.FlushOutbox(context.Background()); err != nil { @@ -115,10 +115,10 @@ func TestCurrentExecuteRejectsInvalidCalleeWithoutStoppingTask(t *testing.T) { } } -func TestCurrentExecuteWaitsForRulesThenDispatchesOriginalIdentity(t *testing.T) { - controller, originator, _, s := newCurrentExecuteFixture(t) - controller.Now = func() time.Time { return currentMonday(8, 59) } - body := currentExecuteBody(t, "waiting-1", "15003164745") +func TestExecuteWaitsForRulesThenDispatchesOriginalIdentity(t *testing.T) { + controller, originator, _, s := newExecuteFixture(t) + controller.Now = func() time.Time { return monday(8, 59) } + body := executeBody(t, "waiting-1", "15003164745") body = []byte(strings.Replace(string(body), `"issued_at":"2026-09-21T01:00:00Z"`, `"issued_at":"2026-09-20T00:00:00Z"`, 1)) if err := controller.ProcessExecute(context.Background(), body); err != nil { t.Fatal(err) @@ -130,16 +130,16 @@ func TestCurrentExecuteWaitsForRulesThenDispatchesOriginalIdentity(t *testing.T) if err != nil || len(pending) != 1 || pending[0].EventID != "waiting-1" { t.Fatalf("pending command was lost: %+v %v", pending, err) } - controller.Now = func() time.Time { return currentMonday(9, 30) } + controller.Now = func() time.Time { return monday(9, 30) } if err := controller.ProcessPending(context.Background()); err != nil || len(originator.calls) != 1 || originator.calls[0].EventID != "waiting-1" { t.Fatalf("original command did not resume once: calls=%+v err=%v", originator.calls, err) } } -func TestCurrentExecuteTimeoutAndMQFailureDoNotRedial(t *testing.T) { - controller, originator, publisher, s := newCurrentExecuteFixture(t) +func TestExecuteTimeoutAndMQFailureDoNotRedial(t *testing.T) { + controller, originator, publisher, s := newExecuteFixture(t) originator.err = context.DeadlineExceeded - body := currentExecuteBody(t, "unknown-1", "15003164745") + body := executeBody(t, "unknown-1", "15003164745") if err := controller.ProcessExecute(context.Background(), body); err == nil { t.Fatal("originator timeout hidden") } @@ -150,7 +150,7 @@ func TestCurrentExecuteTimeoutAndMQFailureDoNotRedial(t *testing.T) { t.Fatalf("unknown call appeared pending: %d %v", len(originator.calls), err) } originator.err = nil - second := currentExecuteBody(t, "mq-loss-1", "15003164745") + second := executeBody(t, "mq-loss-1", "15003164745") if err := controller.ProcessExecute(context.Background(), second); err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_gate_test.go b/internal/dispatcher/gate_test.go similarity index 83% rename from internal/dispatcher/current_gate_test.go rename to internal/dispatcher/gate_test.go index 4bb884e..4d65068 100644 --- a/internal/dispatcher/current_gate_test.go +++ b/internal/dispatcher/gate_test.go @@ -8,8 +8,8 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -func TestCurrentPreDialWindowChangeReleasesUnusedReservation(t *testing.T) { - controller, originator, _, s := newCurrentExecuteFixture(t) +func TestPreDialWindowChangeReleasesUnusedReservation(t *testing.T) { + controller, originator, _, s := newExecuteFixture(t) first := time.Date(2026, 9, 21, 19, 59, 59, 900_000_000, time.FixedZone("Asia/Shanghai", 8*3600)) second := time.Date(2026, 9, 21, 20, 0, 0, 0, time.FixedZone("Asia/Shanghai", 8*3600)) calls := 0 @@ -20,7 +20,7 @@ func TestCurrentPreDialWindowChangeReleasesUnusedReservation(t *testing.T) { } return second } - if err := controller.ProcessExecute(context.Background(), currentExecuteBody(t, "boundary-1", "15003164745")); err != nil { + if err := controller.ProcessExecute(context.Background(), executeBody(t, "boundary-1", "15003164745")); err != nil { t.Fatal(err) } if len(originator.calls) != 0 { @@ -40,10 +40,10 @@ func TestCurrentPreDialWindowChangeReleasesUnusedReservation(t *testing.T) { } } -func TestCurrentFutureIssuedAtNeverDialsEarly(t *testing.T) { - controller, originator, _, s := newCurrentExecuteFixture(t) +func TestFutureIssuedAtNeverDialsEarly(t *testing.T) { + controller, originator, _, s := newExecuteFixture(t) controller.Now = func() time.Time { return time.Date(2026, 9, 21, 8, 59, 59, 0, time.FixedZone("Asia/Shanghai", 8*3600)) } - body := currentExecuteBody(t, "future-1", "15003164745") + body := executeBody(t, "future-1", "15003164745") if err := controller.ProcessExecute(context.Background(), body); err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_policy.go b/internal/dispatcher/policy.go similarity index 62% rename from internal/dispatcher/current_policy.go rename to internal/dispatcher/policy.go index 8e74696..8486800 100644 --- a/internal/dispatcher/current_policy.go +++ b/internal/dispatcher/policy.go @@ -12,12 +12,12 @@ import ( ) var ( - ErrCurrentCalleeRejected = errors.New("callee is not on the approved outbound whitelist") - ErrCurrentRuleWait = errors.New("task rules are temporarily not satisfied") - ErrCurrentRuleInvalid = errors.New("task execution rules are invalid or unknown") + ErrCalleeRejected = errors.New("callee is not on the approved outbound whitelist") + ErrRuleWait = errors.New("task rules are temporarily not satisfied") + ErrRuleInvalid = errors.New("task execution rules are invalid or unknown") ) -type CurrentSelectedTrunk struct { +type SelectedTrunk struct { TrunkID string CallerID string Callee string @@ -26,7 +26,7 @@ type CurrentSelectedTrunk struct { MaxCallDurationMS int64 } -type currentTrunk struct { +type trunkConfig struct { TrunkID string `json:"trunk_id"` Codec string `json:"codec"` DialPrefix string `json:"dial_prefix"` @@ -35,47 +35,47 @@ type currentTrunk struct { AuthMode *string `json:"auth_mode"` RegistrationRequired *bool `json:"registration_required"` MaxConcurrentCalls *int64 `json:"max_concurrent_calls"` - CallerProfiles []currentCallerProfile `json:"caller_profiles"` + CallerProfiles []callerProfile `json:"caller_profiles"` Schedule callwindow.WeeklySchedule `json:"schedule"` } -type currentCallerProfile struct { +type callerProfile struct { ID string `json:"caller_profile_id"` CallerID string `json:"caller_id"` } -var currentCalleeWhitelist = map[string]struct{}{"15003164745": {}, "15830461047": {}} +var calleeWhitelist = map[string]struct{}{"15003164745": {}, "15830461047": {}} -// SelectCurrentTrunk makes one ordered choice before originate. Its answer is +// SelectTrunk makes one ordered choice before originate. Its answer is // frozen with the accepted command; callers never silently reselect on an // unknown execution or after a failed originate. -func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Time, trunkOccupancy, loadedRevisions map[string]int64) (CurrentSelectedTrunk, error) { - if _, allowed := currentCalleeWhitelist[callee]; !allowed { - return CurrentSelectedTrunk{}, ErrCurrentCalleeRejected +func SelectTrunk(snapshot configread.Snapshot, callee string, at time.Time, trunkOccupancy, loadedRevisions map[string]int64) (SelectedTrunk, error) { + if _, allowed := calleeWhitelist[callee]; !allowed { + return SelectedTrunk{}, ErrCalleeRejected } if snapshot.Task.MaxCallDurationMS <= 0 || snapshot.Task.MaxCallDurationMS > math.MaxInt64/int64(time.Millisecond) { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: invalid task call duration", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: invalid task call duration", ErrRuleInvalid) } var taskSchedule callwindow.TaskSchedule if err := json.Unmarshal(snapshot.Task.Schedule, &taskSchedule); err != nil { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: task schedule: %v", ErrCurrentRuleInvalid, err) + return SelectedTrunk{}, fmt.Errorf("%w: task schedule: %v", ErrRuleInvalid, err) } - var trunks []currentTrunk + var trunks []trunkConfig if err := json.Unmarshal(snapshot.SIP.Trunks, &trunks); err != nil { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: SIP trunks: %v", ErrCurrentRuleInvalid, err) + return SelectedTrunk{}, fmt.Errorf("%w: SIP trunks: %v", ErrRuleInvalid, err) } - byID := make(map[string]currentTrunk, len(trunks)) + byID := make(map[string]trunkConfig, len(trunks)) for _, trunk := range trunks { if trunk.TrunkID == "" { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: trunk identity missing", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: trunk identity missing", ErrRuleInvalid) } if _, exists := byID[trunk.TrunkID]; exists { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: duplicate trunk identity", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: duplicate trunk identity", ErrRuleInvalid) } byID[trunk.TrunkID] = trunk } if len(snapshot.Task.AllowedTrunkIDs) == 0 { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: task has no allowed trunks", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: task has no allowed trunks", ErrRuleInvalid) } maxDurationMS := snapshot.Task.MaxCallDurationMS if snapshot.Task.Agent.Mode == "full_ai" { @@ -85,11 +85,11 @@ func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Tim } `json:"conversation"` } if err := json.Unmarshal(snapshot.Task.Agent.Raw, &ai); err != nil { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: AI call duration: %v", ErrCurrentRuleInvalid, err) + return SelectedTrunk{}, fmt.Errorf("%w: AI call duration: %v", ErrRuleInvalid, err) } if ai.Conversation.MaxDurationMS != nil { if *ai.Conversation.MaxDurationMS <= 0 { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: AI call duration must be positive", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: AI call duration must be positive", ErrRuleInvalid) } if *ai.Conversation.MaxDurationMS < maxDurationMS { maxDurationMS = *ai.Conversation.MaxDurationMS @@ -97,17 +97,17 @@ func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Tim } } if maxDurationMS > math.MaxInt64/int64(time.Millisecond) { - return CurrentSelectedTrunk{}, fmt.Errorf("%w: call duration exceeds time range", ErrCurrentRuleInvalid) + return SelectedTrunk{}, fmt.Errorf("%w: call duration exceeds time range", ErrRuleInvalid) } var invalid error for _, id := range snapshot.Task.AllowedTrunkIDs { trunk, found := byID[id] if !found { - invalid = fmt.Errorf("%w: allowed trunk %q is missing", ErrCurrentRuleInvalid, id) + invalid = fmt.Errorf("%w: allowed trunk %q is missing", ErrRuleInvalid, id) continue } if trunk.Transport == nil || *trunk.Transport == "" || trunk.AuthMode == nil || *trunk.AuthMode == "" || trunk.RegistrationRequired == nil || trunk.MaxConcurrentCalls == nil || *trunk.MaxConcurrentCalls <= 0 || trunk.Codec != "PCMA" { - invalid = fmt.Errorf("%w: SIP trunk %q has unknown required execution fields", ErrCurrentRuleInvalid, id) + invalid = fmt.Errorf("%w: SIP trunk %q has unknown required execution fields", ErrRuleInvalid, id) continue } var callerID string @@ -118,7 +118,7 @@ func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Tim } } if callerID == "" { - invalid = fmt.Errorf("%w: caller profile is missing on trunk %q", ErrCurrentRuleInvalid, id) + invalid = fmt.Errorf("%w: caller profile is missing on trunk %q", ErrRuleInvalid, id) continue } end, err := callwindow.Evaluate(taskSchedule, trunk.Schedule, at) @@ -126,7 +126,7 @@ func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Tim if errors.Is(err, callwindow.ErrWindowClosed) { continue } - invalid = fmt.Errorf("%w: schedule of trunk %q: %v", ErrCurrentRuleInvalid, id, err) + invalid = fmt.Errorf("%w: schedule of trunk %q: %v", ErrRuleInvalid, id, err) continue } if !trunk.Enabled || loadedRevisions[id] != snapshot.SIP.Revision || trunkOccupancy[id] >= *trunk.MaxConcurrentCalls { @@ -139,14 +139,14 @@ func SelectCurrentTrunk(snapshot configread.Snapshot, callee string, at time.Tim if !deadline.After(at) { continue } - return CurrentSelectedTrunk{ + return SelectedTrunk{ TrunkID: id, CallerID: callerID, Callee: callee, DialedCallee: trunk.DialPrefix + callee, Deadline: deadline, MaxCallDurationMS: maxDurationMS, }, nil } if invalid != nil { - return CurrentSelectedTrunk{}, invalid + return SelectedTrunk{}, invalid } - return CurrentSelectedTrunk{}, ErrCurrentRuleWait + return SelectedTrunk{}, ErrRuleWait } diff --git a/internal/dispatcher/current_policy_test.go b/internal/dispatcher/policy_test.go similarity index 55% rename from internal/dispatcher/current_policy_test.go rename to internal/dispatcher/policy_test.go index c5937e4..ce6300d 100644 --- a/internal/dispatcher/current_policy_test.go +++ b/internal/dispatcher/policy_test.go @@ -10,26 +10,26 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func currentPolicySnapshot(t *testing.T) configread.Snapshot { +func policySnapshot(t *testing.T) configread.Snapshot { t.Helper() var snapshot configread.Snapshot - if err := json.Unmarshal(currentConfigExample(t, "config-read-task-asr"), &snapshot.Task); err != nil { + if err := json.Unmarshal(configExample(t, "config-read-task-asr"), &snapshot.Task); err != nil { t.Fatal(err) } - sip := string(currentConfigExample(t, "config-read-sip")) + sip := string(configExample(t, "config-read-sip")) for _, replacement := range [][2]string{{`"transport":null`, `"transport":"udp"`}, {`"auth_mode":null`, `"auth_mode":"ip"`}, {`"registration_required":null`, `"registration_required":false`}, {`"max_concurrent_calls":null`, `"max_concurrent_calls":2`}} { sip = strings.Replace(sip, replacement[0], replacement[1], 1) } if err := json.Unmarshal([]byte(sip), &snapshot.SIP); err != nil { t.Fatal(err) } - if err := json.Unmarshal(currentConfigExample(t, "config-read-quota"), &snapshot.Quota); err != nil { + if err := json.Unmarshal(configExample(t, "config-read-quota"), &snapshot.Quota); err != nil { t.Fatal(err) } var providerList struct { Providers []configread.Provider `json:"providers"` } - if err := json.Unmarshal(currentConfigExample(t, "config-read-providers"), &providerList); err != nil { + if err := json.Unmarshal(configExample(t, "config-read-providers"), &providerList); err != nil { t.Fatal(err) } snapshot.Providers = make(map[string]configread.Provider, len(providerList.Providers)) @@ -39,14 +39,14 @@ func currentPolicySnapshot(t *testing.T) configread.Snapshot { return snapshot } -func currentMonday(hour, minute int) time.Time { +func monday(hour, minute int) time.Time { return time.Date(2026, 9, 21, hour, minute, 0, 0, time.FixedZone("Asia/Shanghai", 8*60*60)) } -func TestCurrentSelectTrunkWhitelistScheduleAndCaller(t *testing.T) { - snapshot := currentPolicySnapshot(t) - at := currentMonday(9, 30) - selected, err := SelectCurrentTrunk(snapshot, "15003164745", at, map[string]int64{}, map[string]int64{"trunk-mock": 8}) +func TestSelectTrunkWhitelistScheduleAndCaller(t *testing.T) { + snapshot := policySnapshot(t) + at := monday(9, 30) + selected, err := SelectTrunk(snapshot, "15003164745", at, map[string]int64{}, map[string]int64{"trunk-mock": 8}) if err != nil { t.Fatal(err) } @@ -54,42 +54,42 @@ func TestCurrentSelectTrunkWhitelistScheduleAndCaller(t *testing.T) { t.Fatalf("selected route/prefix/caller/deadline invalid: %+v", selected) } for _, number := range []string{"", "708915003164745", "15003164746", "abc", "158304610470"} { - if _, err := SelectCurrentTrunk(snapshot, number, at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrCurrentCalleeRejected) { + if _, err := SelectTrunk(snapshot, number, at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrCalleeRejected) { t.Fatalf("unapproved number %q not rejected: %v", number, err) } } - for _, at := range []time.Time{currentMonday(8, 59), currentMonday(20, 0)} { - if _, err := SelectCurrentTrunk(snapshot, "15003164745", at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrCurrentRuleWait) { + for _, at := range []time.Time{monday(8, 59), monday(20, 0)} { + if _, err := SelectTrunk(snapshot, "15003164745", at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrRuleWait) { t.Fatalf("outside half-open schedule did not wait: %v", err) } } } -func TestCurrentSelectTrunkFailsClosedForUnknownLineAndLimits(t *testing.T) { - snapshot := currentPolicySnapshot(t) - at := currentMonday(9, 30) - if _, err := SelectCurrentTrunk(snapshot, "15003164745", at, nil, nil); !errors.Is(err, ErrCurrentRuleWait) { +func TestSelectTrunkFailsClosedForUnknownLineAndLimits(t *testing.T) { + snapshot := policySnapshot(t) + at := monday(9, 30) + if _, err := SelectTrunk(snapshot, "15003164745", at, nil, nil); !errors.Is(err, ErrRuleWait) { t.Fatalf("unloaded line did not hold command: %v", err) } - if _, err := SelectCurrentTrunk(snapshot, "15003164745", at, map[string]int64{"trunk-mock": 2}, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrCurrentRuleWait) { + if _, err := SelectTrunk(snapshot, "15003164745", at, map[string]int64{"trunk-mock": 2}, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrRuleWait) { t.Fatalf("full trunk did not hold command: %v", err) } - sip := string(currentConfigExample(t, "config-read-sip")) + sip := string(configExample(t, "config-read-sip")) if err := json.Unmarshal([]byte(sip), &snapshot.SIP); err != nil { t.Fatal(err) } - if _, err := SelectCurrentTrunk(snapshot, "15003164745", at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrCurrentRuleInvalid) { + if _, err := SelectTrunk(snapshot, "15003164745", at, nil, map[string]int64{"trunk-mock": 8}); !errors.Is(err, ErrRuleInvalid) { t.Fatalf("unknown transport/auth/quota was not fail-closed: %v", err) } } -func TestCurrentSelectTrunkCapsCallByAIAndPreservesPrefix(t *testing.T) { - snapshot := currentPolicySnapshot(t) +func TestSelectTrunkCapsCallByAIAndPreservesPrefix(t *testing.T) { + snapshot := policySnapshot(t) snapshot.Task.Agent.Mode = "full_ai" snapshot.Task.Agent.Raw = []byte(`{"mode":"full_ai","conversation":{"max_duration_ms":60000}}`) snapshot.SIP.Trunks = []byte(strings.Replace(string(snapshot.SIP.Trunks), `"dial_prefix":""`, `"dial_prefix":"7089"`, 1)) - at := currentMonday(9, 30) - selected, err := SelectCurrentTrunk(snapshot, "15830461047", at, nil, map[string]int64{"trunk-mock": 8}) + at := monday(9, 30) + selected, err := SelectTrunk(snapshot, "15830461047", at, nil, map[string]int64{"trunk-mock": 8}) if err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_runtime.go b/internal/dispatcher/runtime.go similarity index 88% rename from internal/dispatcher/current_runtime.go rename to internal/dispatcher/runtime.go index b9a2caf..e50dcf9 100644 --- a/internal/dispatcher/current_runtime.go +++ b/internal/dispatcher/runtime.go @@ -15,13 +15,13 @@ import ( "git.ipao.vip/rogee/go-sip/internal/tenant" ) -// CurrentRuntime owns task/control consumers for one Dispatcher. SaaS creates +// Runtime owns task/control consumers for one Dispatcher. SaaS creates // every queue/binding; this process only checks and consumes predeclared ones. -type CurrentRuntime struct { +type Runtime struct { Broker *mq.Broker - Bootstrap CurrentBootstrap - Execute CurrentExecuteController - Control CurrentControlController + Bootstrap Bootstrap + Execute ExecuteController + Control ControlController PollInterval time.Duration DiscoveryInterval time.Duration Logger *slog.Logger @@ -34,9 +34,9 @@ type CurrentRuntime struct { // Serve closes admission on every shutdown/failure; it never clears durable // calls, results, task queues, or the SQLite file. -func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { +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.DiscoveryInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil { - return errors.New("current runtime requires one Dispatcher, durable state, verified SIP, independent clocks, and configured polling; external control drain is forbidden") + return errors.New("runtime requires one Dispatcher, durable state, verified SIP, independent clocks, and configured polling; external control drain is forbidden") } if err := r.Execute.validate(); err != nil { return err @@ -93,7 +93,7 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { } r.Logger.Warn("control and outbox stay active while newer SIP revision waits; task admission remains closed", "dispatcher_id", r.Bootstrap.DispatcherID, "error", err) } - follower := &CurrentDiscoveryFollower{DispatcherID: r.Bootstrap.DispatcherID, Client: r.Bootstrap.Client, Store: r.Bootstrap.Store, ApprovedSIP: sip, VerifySIP: r.Bootstrap.VerifySIP, Cursor: cursor} + follower := &DiscoveryFollower{DispatcherID: r.Bootstrap.DispatcherID, Client: r.Bootstrap.Client, Store: r.Bootstrap.Store, ApprovedSIP: sip, VerifySIP: r.Bootstrap.VerifySIP, Cursor: cursor} if err := r.syncTaskConsumers(ctx, consumers); err != nil { return err } @@ -140,7 +140,7 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { } } -func (r *CurrentRuntime) watchConsumer(ctx context.Context, queue string, consumer *mq.Consumer) { +func (r *Runtime) watchConsumer(ctx context.Context, queue string, consumer *mq.Consumer) { err := consumer.Wait(ctx) if err != nil && ctx.Err() == nil { r.Logger.Error("current MQ consumer failed", "dispatcher_id", r.Bootstrap.DispatcherID, "queue", queue, "error", err) @@ -153,7 +153,7 @@ func (r *CurrentRuntime) watchConsumer(ctx context.Context, queue string, consum // signalFailure forces the runtime to close admission after a consumer // handler fails; RabbitMQ requeue alone would otherwise spin indefinitely. -func (r *CurrentRuntime) signalFailure(err error) { +func (r *Runtime) signalFailure(err error) { if r.failures != nil { select { case r.failures <- err: @@ -162,7 +162,7 @@ func (r *CurrentRuntime) signalFailure(err error) { } } -func (r *CurrentRuntime) withTask(taskID string, process func() error) error { +func (r *Runtime) withTask(taskID string, process func() error) error { r.gate.RLock() defer r.gate.RUnlock() r.locksMu.Lock() @@ -177,7 +177,7 @@ func (r *CurrentRuntime) withTask(taskID string, process func() error) error { return process() } -func (r *CurrentRuntime) handleControl(ctx context.Context, _ string, body []byte) error { +func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) error { var event struct { EventType string `json:"event_type"` Payload struct { @@ -218,7 +218,7 @@ func (r *CurrentRuntime) handleControl(ctx context.Context, _ string, body []byt } } -func (r *CurrentRuntime) handleTask(ctx context.Context, taskID string, body []byte) error { +func (r *Runtime) handleTask(ctx context.Context, taskID string, body []byte) error { return r.withTask(taskID, func() error { if err := r.Execute.ProcessExecute(ctx, body); err != nil { r.Logger.Error("call instruction failed", "dispatcher_id", r.Bootstrap.DispatcherID, "task_id", taskID, "error", err) @@ -229,7 +229,7 @@ func (r *CurrentRuntime) handleTask(ctx context.Context, taskID string, body []b }) } -func (r *CurrentRuntime) processPending(ctx context.Context) error { +func (r *Runtime) processPending(ctx context.Context) error { pending, err := r.Bootstrap.Store.ListPendingExecute(r.Bootstrap.DispatcherID) if err != nil { return fmt.Errorf("read durable rule-wait instructions: %w", err) @@ -243,7 +243,7 @@ func (r *CurrentRuntime) processPending(ctx context.Context) error { return nil } -func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error { +func (r *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error { assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID) if err != nil { return err @@ -320,4 +320,4 @@ func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[st // Compile-time interface checks: the same broker confirms bound persistent // results and consumes SaaS-owned queues without configure permissions. -var _ CurrentPublisher = (*mq.Broker)(nil) +var _ Publisher = (*mq.Broker)(nil) diff --git a/internal/dispatcher/current_runtime_integration_test.go b/internal/dispatcher/runtime_integration_test.go similarity index 81% rename from internal/dispatcher/current_runtime_integration_test.go rename to internal/dispatcher/runtime_integration_test.go index 0e2bf70..85ff918 100644 --- a/internal/dispatcher/current_runtime_integration_test.go +++ b/internal/dispatcher/runtime_integration_test.go @@ -24,24 +24,24 @@ import ( amqp "github.com/rabbitmq/amqp091-go" ) -type currentRuntimeMockAgent struct { - calls chan CurrentCallSpec - controls chan CurrentControlSpec +type runtimeMockAgent struct { + calls chan CallSpec + controls chan ControlSpec } -func (a *currentRuntimeMockAgent) LoadedTrunks(context.Context) (map[string]int64, error) { +func (a *runtimeMockAgent) LoadedTrunks(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-mock": 8}, nil } -func (a *currentRuntimeMockAgent) Originate(_ context.Context, spec CurrentCallSpec) error { +func (a *runtimeMockAgent) Originate(_ context.Context, spec CallSpec) error { a.calls <- spec return nil } -func (a *currentRuntimeMockAgent) SendControl(_ context.Context, spec CurrentControlSpec) error { +func (a *runtimeMockAgent) SendControl(_ context.Context, spec ControlSpec) error { a.controls <- spec return nil } -func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { +func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { brokerURL, adminURL := os.Getenv("RABBITMQ_URL"), os.Getenv("RABBITMQ_PROVISIONER_URL") if brokerURL == "" || adminURL == "" { t.Skip("requires isolated RabbitMQ mock with provisioner account") @@ -91,10 +91,10 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T t.Fatal(err) } } - publish(controlRoute, currentControlBody(t, "control-example", "pause", "drain")) - publish(taskRoute, currentExecuteBody(t, "call-example", "15003164745")) + publish(controlRoute, controlBody(t, "control-example", "pause", "drain")) + publish(taskRoute, executeBody(t, "call-example", "15003164745")) - snapshot := currentPolicySnapshot(t) + snapshot := policySnapshot(t) sipJSON, err := json.Marshal(snapshot.SIP) if err != nil { t.Fatal(err) @@ -107,19 +107,19 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T body = sipJSON case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - body = currentConfigExample(t, "task-discovery-page") + body = configExample(t, "task-discovery-page") } else if r.URL.Query().Get("after") != "" { - body = currentConfigExample(t, "task-discovery-end") + body = configExample(t, "task-discovery-end") } else { w.WriteHeader(http.StatusBadRequest) return } case "/internal/v1/dispatcher/task/task-asr": - body = currentConfigExample(t, "config-read-task-asr") + body = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/ai-providers": - body = currentConfigExample(t, "config-read-providers") + body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/tenant/1001/quota": - body = currentConfigExample(t, "config-read-quota") + body = configExample(t, "config-read-quota") default: t.Errorf("unexpected HTTP configuration path %s", r.URL.Path) w.WriteHeader(http.StatusNotFound) @@ -142,7 +142,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T t.Fatal(err) } defer broker.Close() - agent := ¤tRuntimeMockAgent{calls: make(chan CurrentCallSpec, 3), controls: make(chan CurrentControlSpec, 3)} + agent := &runtimeMockAgent{calls: make(chan CallSpec, 3), controls: make(chan ControlSpec, 3)} var windowAllowed atomic.Bool windowAllowed.Store(true) verify := func(_ context.Context, sip configread.SIP) error { @@ -151,16 +151,16 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T } return nil } - runtime := &CurrentRuntime{ + runtime := &Runtime{ Broker: broker, - Bootstrap: CurrentBootstrap{DispatcherID: id, Client: client, Store: db, VerifySIP: verify}, - Execute: CurrentExecuteController{DispatcherID: id, Store: db, Originator: agent, Publisher: broker, Now: func() time.Time { + Bootstrap: Bootstrap{DispatcherID: id, Client: client, Store: db, VerifySIP: verify}, + Execute: ExecuteController{DispatcherID: id, Store: db, Originator: agent, Publisher: broker, Now: func() time.Time { if windowAllowed.Load() { - return currentMonday(9, 30) + return monday(9, 30) } - return currentMonday(8, 59) + return monday(8, 59) }}, - Control: CurrentControlController{DispatcherID: id, Store: db, Client: client, Agent: agent, VerifySIP: verify, Now: func() time.Time { return currentMonday(9, 30) }}, + Control: ControlController{DispatcherID: id, Store: db, Client: client, Agent: agent, VerifySIP: verify, Now: func() time.Time { return monday(9, 30) }}, PollInterval: 30 * time.Millisecond, DiscoveryInterval: 120 * time.Millisecond, Logger: slog.Default(), } ctx, cancel := context.WithCancel(context.Background()) @@ -198,7 +198,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T if err != nil || state.Messages != 1 { t.Fatalf("task queue was consumed before resume: %+v %v", state, err) } - publish(controlRoute, currentControlBody(t, "resume-integration", "resume", "")) + publish(controlRoute, controlBody(t, "resume-integration", "resume", "")) select { case spec := <-agent.controls: if spec.Action != "resume" { @@ -253,7 +253,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T if seen["control-example"] != "applied" || seen["resume-integration"] != "applied" || seen["call-example"] != "dispatched" { t.Fatalf("wrong shared SaaS results: %+v", seen) } - publish(taskRoute, currentExecuteBody(t, "call-example", "15003164745")) + publish(taskRoute, executeBody(t, "call-example", "15003164745")) select { case spec := <-agent.calls: t.Fatalf("redelivery reoriginated call %s", spec.EventID) @@ -263,7 +263,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T // stops so further instructions remain in SaaS's task queue. Admission // resumes from the original identity when the configured window opens. windowAllowed.Store(false) - waiting := []byte(strings.Replace(string(currentExecuteBody(t, "window-wait-1", "15003164745")), `"issued_at":"2026-09-21T01:00:00Z"`, `"issued_at":"2026-09-20T00:00:00Z"`, 1)) + waiting := []byte(strings.Replace(string(executeBody(t, "window-wait-1", "15003164745")), `"issued_at":"2026-09-21T01:00:00Z"`, `"issued_at":"2026-09-20T00:00:00Z"`, 1)) publish(taskRoute, waiting) waitDeadline := time.After(5 * time.Second) for { @@ -284,7 +284,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T case <-time.After(20 * time.Millisecond): } } - publish(taskRoute, currentExecuteBody(t, "window-wait-2", "15003164745")) + publish(taskRoute, executeBody(t, "window-wait-2", "15003164745")) time.Sleep(100 * time.Millisecond) queueState, err := admin.QueueInspect(taskRoute.Queue) if err != nil || queueState.Messages != 1 { @@ -309,7 +309,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T // A SIP notification is ACKed only after its revision is durable. The // already-dispatched calls keep the revision change fenced, while the // independent control queue and shared result publisher continue. - sipChange := []byte(strings.Replace(string(currentConfigExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) + sipChange := []byte(strings.Replace(string(configExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) publish(controlRoute, sipChange) deadline := time.After(5 * time.Second) for { @@ -326,13 +326,13 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T case <-time.After(20 * time.Millisecond): } } - publish(taskRoute, currentExecuteBody(t, "after-sip-change", "15003164745")) + publish(taskRoute, executeBody(t, "after-sip-change", "15003164745")) select { case spec := <-agent.calls: t.Fatalf("pending SIP change originated call %s", spec.EventID) case <-time.After(120 * time.Millisecond): } - publish(controlRoute, currentControlBody(t, "stop-after-sip", "stop", "")) + publish(controlRoute, controlBody(t, "stop-after-sip", "stop", "")) select { case spec := <-agent.controls: if spec.Action != "stop" || spec.ActiveCallPolicy != "hangup" { diff --git a/internal/dispatcher/current_sip_conflict_test.go b/internal/dispatcher/sip_conflict_test.go similarity index 76% rename from internal/dispatcher/current_sip_conflict_test.go rename to internal/dispatcher/sip_conflict_test.go index ebb0838..099d888 100644 --- a/internal/dispatcher/current_sip_conflict_test.go +++ b/internal/dispatcher/sip_conflict_test.go @@ -12,7 +12,7 @@ import ( "git.ipao.vip/rogee/go-sip/internal/store" ) -func TestCurrentBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) { +func TestBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" sipReads := 0 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -21,22 +21,22 @@ func TestCurrentBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) switch r.URL.Path { case "/internal/v1/dispatcher/sip": sipReads++ - response = currentConfigExample(t, "config-read-sip") + response = configExample(t, "config-read-sip") if sipReads > 1 { response = []byte(strings.Replace(string(response), `"dial_prefix":""`, `"dial_prefix":"9"`, 1)) } case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - response = currentConfigExample(t, "task-discovery-page") + response = configExample(t, "task-discovery-page") } else { - response = currentConfigExample(t, "task-discovery-end") + response = configExample(t, "task-discovery-end") } case "/internal/v1/dispatcher/ai-providers": - response = currentConfigExample(t, "config-read-providers") + response = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - response = currentConfigExample(t, "config-read-task-asr") + response = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - response = currentConfigExample(t, "config-read-quota") + response = configExample(t, "config-read-quota") default: t.Errorf("unexpected path %s", r.URL.Path) } @@ -52,7 +52,7 @@ func TestCurrentBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) t.Fatal(err) } defer db.Close() - bootstrap := CurrentBootstrap{Client: client, Store: db, DispatcherID: id, + bootstrap := Bootstrap{Client: client, Store: db, DispatcherID: id, VerifySIP: func(context.Context, configread.SIP) error { return nil }, DrainControls: func(context.Context) error { return nil }, } diff --git a/internal/dispatcher/current_sip_reload.go b/internal/dispatcher/sip_reload.go similarity index 95% rename from internal/dispatcher/current_sip_reload.go rename to internal/dispatcher/sip_reload.go index 8ab04f0..9db75ca 100644 --- a/internal/dispatcher/current_sip_reload.go +++ b/internal/dispatcher/sip_reload.go @@ -11,7 +11,7 @@ import ( // 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. -func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscoveryFollower) error { +func (r *Runtime) refreshSIP(ctx context.Context, follower *DiscoveryFollower) error { if r == nil || follower == 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") } @@ -93,7 +93,7 @@ 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 { + if err := validateAISnapshot(snapshot); err != nil { return r.closeSIPAdmission(err) } if err := r.Bootstrap.Store.SaveSnapshot(snapshot); err != nil { @@ -122,7 +122,7 @@ func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscov return nil } -func (r *CurrentRuntime) closeSIPAdmission(cause error) error { +func (r *Runtime) closeSIPAdmission(cause error) error { if err := r.Bootstrap.Store.CloseAdmission(r.Bootstrap.DispatcherID); err != nil { return errors.Join(cause, fmt.Errorf("close admission after SIP refresh failure: %w", err)) } diff --git a/internal/dispatcher/current_sip_runtime_test.go b/internal/dispatcher/sip_runtime_test.go similarity index 73% rename from internal/dispatcher/current_sip_runtime_test.go rename to internal/dispatcher/sip_runtime_test.go index a4be074..8229f75 100644 --- a/internal/dispatcher/current_sip_runtime_test.go +++ b/internal/dispatcher/sip_runtime_test.go @@ -14,9 +14,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func TestCurrentSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing.T) { - executor, _, _, s := newCurrentExecuteFixture(t) - approved := currentPolicySnapshot(t).SIP +func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing.T) { + executor, _, _, s := newExecuteFixture(t) + approved := policySnapshot(t).SIP changed := approved changed.Revision = 9 newSIP, err := json.Marshal(changed) @@ -31,16 +31,16 @@ func TestCurrentSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *t body = newSIP case "/internal/v1/dispatcher/tasks": if r.URL.Query().Get("after") == "" { - body = currentConfigExample(t, "task-discovery-page") + body = configExample(t, "task-discovery-page") } else { - body = currentConfigExample(t, "task-discovery-end") + body = configExample(t, "task-discovery-end") } case "/internal/v1/dispatcher/ai-providers": - body = currentConfigExample(t, "config-read-providers") + body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/task/task-asr": - body = currentConfigExample(t, "config-read-task-asr") + body = configExample(t, "config-read-task-asr") case "/internal/v1/dispatcher/tenant/1001/quota": - body = currentConfigExample(t, "config-read-quota") + body = configExample(t, "config-read-quota") default: t.Errorf("unexpected HTTP path %s", r.URL.Path) w.WriteHeader(http.StatusNotFound) @@ -60,9 +60,9 @@ func TestCurrentSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *t } return nil } - runtime := &CurrentRuntime{Bootstrap: CurrentBootstrap{DispatcherID: executor.DispatcherID, Client: client, Store: s, VerifySIP: verify}, Logger: slog.Default()} - follower := &CurrentDiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: verify, Cursor: "opaque-end-token"} - body := []byte(strings.Replace(string(currentConfigExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) + runtime := &Runtime{Bootstrap: Bootstrap{DispatcherID: executor.DispatcherID, Client: client, Store: s, VerifySIP: verify}, Logger: slog.Default()} + follower := &DiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: verify, Cursor: "opaque-end-token"} + body := []byte(strings.Replace(string(configExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) if err := runtime.handleControl(context.Background(), "", body); err != nil { t.Fatal(err) } diff --git a/internal/rpc/approved_control_transport_test.go b/internal/rpc/approved_control_transport_test.go index 6fa0674..423a640 100644 --- a/internal/rpc/approved_control_transport_test.go +++ b/internal/rpc/approved_control_transport_test.go @@ -68,7 +68,7 @@ func TestApprovedTaskControlOverMutualTLSRequiresPinnedPeerAndSession(t *testing defer release() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() - spec := dispatcher.CurrentControlSpec{DispatcherID: task.DispatcherID, TenantID: task.TenantID, TaskID: task.TaskID, Action: "pause", ActiveCallPolicy: "hangup"} + spec := dispatcher.ControlSpec{DispatcherID: task.DispatcherID, TenantID: task.TenantID, TaskID: task.TaskID, Action: "pause", ActiveCallPolicy: "hangup"} completed := make(chan error, 1) go func() { completed <- originator.SendControl(ctx, spec) }() select { diff --git a/internal/rpc/approved_full_ai_integration_test.go b/internal/rpc/approved_full_ai_integration_test.go index 5085c91..365d597 100644 --- a/internal/rpc/approved_full_ai_integration_test.go +++ b/internal/rpc/approved_full_ai_integration_test.go @@ -133,7 +133,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if err := orig.VerifySIP(context.Background(), sip); err != nil { t.Fatal(err) } - spec := dispatcher.CurrentCallSpec{ + spec := dispatcher.CallSpec{ EventID: req.SourceEventId, TaskID: req.TaskId, TenantID: req.TenantId, TrunkID: req.SelectedTrunkId, Callee: req.Callee, DialedCallee: req.DialedCallee, CallerID: req.CallerId, RingTimeoutMS: req.RingTimeoutMs, diff --git a/internal/rpc/approved_integration_test.go b/internal/rpc/approved_integration_test.go index 3b78398..6aa66dd 100644 --- a/internal/rpc/approved_integration_test.go +++ b/internal/rpc/approved_integration_test.go @@ -72,7 +72,7 @@ func TestApprovedDispatcherToAgentUnaryMockRetainsSnapshotAndOneShotCall(t *test if err := orig.VerifySIP(context.Background(), sip); err != nil { t.Fatal(err) } - spec := dispatcher.CurrentCallSpec{ + spec := dispatcher.CallSpec{ EventID: req.SourceEventId, TaskID: task.TaskID, TenantID: task.TenantID, TrunkID: req.SelectedTrunkId, Callee: req.Callee, DialedCallee: req.DialedCallee, CallerID: req.CallerId, RingTimeoutMS: req.RingTimeoutMs,