From 16c7ca24970bec585f122f95da6038ea3310ac68 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 11:56:58 +0800 Subject: [PATCH] Normalize configuration reader API and paths --- .../current_dispatcher_setup_test.go | 2 +- .../saas-dispatcher-implementation.md | 1 + internal/ai/current.go | 16 ++-- internal/ai/current_test.go | 24 ++--- .../{current_discovery.go => discovery.go} | 38 ++++---- ...nt_discovery_test.go => discovery_test.go} | 16 ++-- .../configread/{current.go => snapshots.go} | 92 +++++++++---------- .../{current_test.go => snapshots_test.go} | 50 +++++----- internal/dispatcher/approved_originator.go | 2 +- .../dispatcher/approved_originator_test.go | 2 +- internal/dispatcher/current_ai.go | 2 +- internal/dispatcher/current_ai_test.go | 4 +- internal/dispatcher/current_config.go | 12 +-- internal/dispatcher/current_config_test.go | 6 +- internal/dispatcher/current_control.go | 4 +- internal/dispatcher/current_control_test.go | 6 +- internal/dispatcher/current_discovery.go | 8 +- internal/dispatcher/current_discovery_test.go | 2 +- internal/dispatcher/current_execute.go | 2 +- internal/dispatcher/current_execute_test.go | 2 +- internal/dispatcher/current_policy.go | 2 +- internal/dispatcher/current_policy_test.go | 8 +- internal/dispatcher/current_runtime.go | 2 +- .../current_runtime_integration_test.go | 2 +- .../dispatcher/current_sip_conflict_test.go | 2 +- internal/dispatcher/current_sip_reload.go | 6 +- .../dispatcher/current_sip_runtime_test.go | 2 +- internal/rpc/approved_execution.go | 4 +- internal/rpc/approved_execution_test.go | 6 +- .../rpc/approved_full_ai_integration_test.go | 8 +- internal/rpc/approved_integration_test.go | 8 +- internal/rpc/approved_runner_test.go | 8 +- internal/rpc/recording_server_flow_test.go | 12 +-- internal/store/calls.go | 6 +- internal/store/calls_test.go | 2 +- internal/store/discovery_test.go | 20 ++-- internal/store/global_revision_test.go | 2 +- internal/store/result.go | 2 +- internal/store/store.go | 44 ++++----- internal/store/store_test.go | 14 +-- 40 files changed, 226 insertions(+), 225 deletions(-) rename internal/configread/{current_discovery.go => discovery.go} (59%) rename internal/configread/{current_discovery_test.go => discovery_test.go} (87%) rename internal/configread/{current.go => snapshots.go} (59%) rename internal/configread/{current_test.go => snapshots_test.go} (74%) diff --git a/cmd/sip-go-agent/current_dispatcher_setup_test.go b/cmd/sip-go-agent/current_dispatcher_setup_test.go index f23f95f..7f2b54f 100644 --- a/cmd/sip-go-agent/current_dispatcher_setup_test.go +++ b/cmd/sip-go-agent/current_dispatcher_setup_test.go @@ -44,7 +44,7 @@ func TestDispatcherConfigurationClientReadsApprovedSIPWithExplicitIdentity(t *te if settings.DispatcherID != dispatcherID || settings.SQLitePath != path { t.Fatal("dispatcher startup changed approved deployment identity or SQLite path") } - sip, err := reader.ReadCurrentSIP(context.Background()) + sip, err := reader.ReadSIP(context.Background()) if err != nil { t.Fatal(err) } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 0112ecd..89b93ef 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -108,6 +108,7 @@ - Store 名称收敛:移除现行 Store 的自有 `Current*` 类型、错误和 `OpenCurrent` 入口,保留单一 `Store`/`Open`;20 份 Go 源码与测试文件改为不带代次的路径,原现行任务/录音/结果测试仍执行。只重命名源码与调用,不修改现有 SQLite 表、记录或恢复数据。 - 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 路径及断连拒绝行为未改变;相关调用方和隔离测试同步更新。 ## 验收台账 diff --git a/internal/ai/current.go b/internal/ai/current.go index 647b6d4..0d3e302 100644 --- a/internal/ai/current.go +++ b/internal/ai/current.go @@ -30,13 +30,13 @@ type CurrentBound struct { } type CurrentASR struct { - Provider configread.CurrentProvider + Provider configread.Provider Request doubaospeech.ASRV2Config Timeout time.Duration } type CurrentLLM struct { - Provider configread.CurrentProvider + Provider configread.Provider Model string Temperature *float64 MaxTokens *int64 @@ -44,7 +44,7 @@ type CurrentLLM struct { } type CurrentTTS struct { - Provider configread.CurrentProvider + Provider configread.Provider Request doubaospeech.TTSV2Request Timeout time.Duration } @@ -114,14 +114,14 @@ type currentAgentSettings struct { // BindCurrent rejects schema-valid settings which the selected SDK cannot // express. In particular, the published TTS schema is intentionally not // silently narrowed to the SDK's speed/format capabilities. -func BindCurrent(task configread.CurrentTask, providers map[string]configread.CurrentProvider) (CurrentBound, error) { +func BindCurrent(task configread.Task, providers map[string]configread.Provider) (CurrentBound, error) { if len(task.Raw) == 0 { return CurrentBound{}, errors.New("approved task snapshot is missing") } if err := contract.ValidateCurrent("config-read", task.Raw); err != nil { return CurrentBound{}, fmt.Errorf("task snapshot: %w", err) } - var frozen configread.CurrentTask + var frozen configread.Task if err := json.Unmarshal(task.Raw, &frozen); err != nil { return CurrentBound{}, fmt.Errorf("decode task snapshot: %w", err) } @@ -254,14 +254,14 @@ func BindCurrent(task configread.CurrentTask, providers map[string]configread.Cu return bound, nil } -func currentProvider(providers map[string]configread.CurrentProvider, ref, role, adapter string) (configread.CurrentProvider, error) { +func currentProvider(providers map[string]configread.Provider, ref, role, adapter string) (configread.Provider, error) { p, found := providers[ref] if !found || ref == "" || p.ProviderRef != ref || !p.Enabled || p.Role != role || p.Adapter != adapter || p.Credential == "" { - return configread.CurrentProvider{}, fmt.Errorf("%s provider is missing, disabled, or incompatible", role) + return configread.Provider{}, fmt.Errorf("%s provider is missing, disabled, or incompatible", role) } u, err := url.Parse(p.Endpoint) if err != nil || u.Host == "" || u.User != nil || (u.Scheme != "https" && u.Scheme != "http" && u.Scheme != "wss" && u.Scheme != "ws") { - return configread.CurrentProvider{}, fmt.Errorf("%s provider endpoint is invalid", role) + return configread.Provider{}, fmt.Errorf("%s provider endpoint is invalid", role) } return p, nil } diff --git a/internal/ai/current_test.go b/internal/ai/current_test.go index 814551f..cb3bfb6 100644 --- a/internal/ai/current_test.go +++ b/internal/ai/current_test.go @@ -11,7 +11,7 @@ import ( doubaospeech "github.com/GizClaw/doubao-speech-go" ) -func currentFixture(t *testing.T, mode string) (configread.CurrentTask, map[string]configread.CurrentProvider) { +func currentFixture(t *testing.T, mode string) (configread.Task, map[string]configread.Provider) { t.Helper() name := "config-read-task-full.json" if mode == "asr_only" { @@ -21,7 +21,7 @@ func currentFixture(t *testing.T, mode string) (configread.CurrentTask, map[stri if err != nil { t.Fatal(err) } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(raw, &task); err != nil { t.Fatal(err) } @@ -30,19 +30,19 @@ func currentFixture(t *testing.T, mode string) (configread.CurrentTask, map[stri t.Fatal(err) } var list struct { - Providers []configread.CurrentProvider `json:"providers"` + Providers []configread.Provider `json:"providers"` } if err := json.Unmarshal(raw, &list); err != nil { t.Fatal(err) } - providers := make(map[string]configread.CurrentProvider, len(list.Providers)) + providers := make(map[string]configread.Provider, len(list.Providers)) for _, p := range list.Providers { providers[p.ProviderRef] = p } return task, providers } -func changeCurrentAgent(t *testing.T, task configread.CurrentTask, change func(map[string]any)) configread.CurrentTask { +func changeCurrentAgent(t *testing.T, task configread.Task, change func(map[string]any)) configread.Task { t.Helper() var body map[string]any if err := json.Unmarshal(task.Raw, &body); err != nil { @@ -54,7 +54,7 @@ func changeCurrentAgent(t *testing.T, task configread.CurrentTask, change func(m if err != nil { t.Fatal(err) } - var updated configread.CurrentTask + var updated configread.Task if err := json.Unmarshal(raw, &updated); err != nil { t.Fatal(err) } @@ -122,25 +122,25 @@ func TestBindCurrentRejectsSDKUnsupportedTTSWithoutChangingSchema(t *testing.T) func TestBindCurrentRejectsUnauthorizedProviderBeforeCall(t *testing.T) { for _, tc := range []struct { name string - mutate func(map[string]configread.CurrentProvider) + mutate func(map[string]configread.Provider) }{ - {"missing", func(ps map[string]configread.CurrentProvider) { delete(ps, "asr-example") }}, - {"disabled", func(ps map[string]configread.CurrentProvider) { + {"missing", func(ps map[string]configread.Provider) { delete(ps, "asr-example") }}, + {"disabled", func(ps map[string]configread.Provider) { p := ps["asr-example"] p.Enabled = false ps[p.ProviderRef] = p }}, - {"wrong-role", func(ps map[string]configread.CurrentProvider) { + {"wrong-role", func(ps map[string]configread.Provider) { p := ps["asr-example"] p.Role = "tts" ps[p.ProviderRef] = p }}, - {"wrong-adapter", func(ps map[string]configread.CurrentProvider) { + {"wrong-adapter", func(ps map[string]configread.Provider) { p := ps["asr-example"] p.Adapter = "unknown" ps[p.ProviderRef] = p }}, - {"missing-credential", func(ps map[string]configread.CurrentProvider) { + {"missing-credential", func(ps map[string]configread.Provider) { p := ps["asr-example"] p.Credential = "" ps[p.ProviderRef] = p diff --git a/internal/configread/current_discovery.go b/internal/configread/discovery.go similarity index 59% rename from internal/configread/current_discovery.go rename to internal/configread/discovery.go index 94a3232..4f06515 100644 --- a/internal/configread/current_discovery.go +++ b/internal/configread/discovery.go @@ -10,64 +10,64 @@ import ( "git.ipao.vip/rogee/go-sip/internal/contract" ) -// CurrentTaskPage is one page of the current, cursor-based discovery list. +// TaskPage is one page of the current, cursor-based discovery list. // A page ends the list only when Tasks is empty; short pages do not end it. -type CurrentTaskPage struct { - DispatcherID string `json:"dispatcher_id"` - Cursor string `json:"cursor"` - Tasks []CurrentDiscoveredTask `json:"tasks"` +type TaskPage struct { + DispatcherID string `json:"dispatcher_id"` + Cursor string `json:"cursor"` + Tasks []DiscoveredTask `json:"tasks"` } -type CurrentDiscoveredTask struct { +type DiscoveredTask struct { TaskID string `json:"task_id"` TenantID int64 `json:"tenant_id"` Status string `json:"status"` TaskRevision int64 `json:"task_revision"` } -// ReadCurrentTasks retrieves exactly one page, with no old snapshot/watermark +// ReadTasks retrieves exactly one page, with no old snapshot/watermark // fallback. The caller must persist non-empty pages before advancing after. -func (c *Client) ReadCurrentTasks(ctx context.Context, after string) (CurrentTaskPage, error) { +func (c *Client) ReadTasks(ctx context.Context, after string) (TaskPage, error) { path := configReadPath + "/tasks" if after != "" { path += "?after=" + url.QueryEscape(after) } body, err := c.get(ctx, path) if err != nil { - return CurrentTaskPage{}, fmt.Errorf("read task discovery: %w", err) + return TaskPage{}, fmt.Errorf("read task discovery: %w", err) } if err := contract.ValidateCurrent("task-discovery", body); err != nil { - return CurrentTaskPage{}, err + return TaskPage{}, err } - var page CurrentTaskPage + var page TaskPage if err := json.Unmarshal(body, &page); err != nil { - return CurrentTaskPage{}, fmt.Errorf("decode task discovery: %w", err) + return TaskPage{}, fmt.Errorf("decode task discovery: %w", err) } if page.DispatcherID != c.dispatcherID { - return CurrentTaskPage{}, errors.New("task discovery dispatcher owner mismatch") + return TaskPage{}, errors.New("task discovery dispatcher owner mismatch") } if len(page.Tasks) > 0 && page.Cursor == after { - return CurrentTaskPage{}, errors.New("non-empty task discovery page did not advance cursor") + return TaskPage{}, errors.New("non-empty task discovery page did not advance cursor") } seen := make(map[string]struct{}, len(page.Tasks)) for _, task := range page.Tasks { if _, exists := seen[task.TaskID]; exists { - return CurrentTaskPage{}, fmt.Errorf("duplicate task %q in discovery page", task.TaskID) + return TaskPage{}, fmt.Errorf("duplicate task %q in discovery page", task.TaskID) } seen[task.TaskID] = struct{}{} } return page, nil } -// ReadAllCurrentTasks collects a cold-start list. Callers must validate and +// ReadAllTasks collects a cold-start list. Callers must validate and // commit it atomically before opening task admission; no page is durable here. -func (c *Client) ReadAllCurrentTasks(ctx context.Context) ([]CurrentDiscoveredTask, string, error) { - var tasks []CurrentDiscoveredTask +func (c *Client) ReadAllTasks(ctx context.Context) ([]DiscoveredTask, string, error) { + var tasks []DiscoveredTask seenTasks := make(map[string]struct{}) seenCursors := make(map[string]struct{}) cursor := "" for { - page, err := c.ReadCurrentTasks(ctx, cursor) + page, err := c.ReadTasks(ctx, cursor) if err != nil { return nil, "", err } diff --git a/internal/configread/current_discovery_test.go b/internal/configread/discovery_test.go similarity index 87% rename from internal/configread/current_discovery_test.go rename to internal/configread/discovery_test.go index 4512b0b..ebbdd90 100644 --- a/internal/configread/current_discovery_test.go +++ b/internal/configread/discovery_test.go @@ -10,7 +10,7 @@ import ( "testing" ) -func TestReadAllCurrentTasksContinuesAfterShortPageUntilEmptyPage(t *testing.T) { +func TestReadAllTasksContinuesAfterShortPageUntilEmptyPage(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" queries := []string{} server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -35,7 +35,7 @@ func TestReadAllCurrentTasksContinuesAfterShortPageUntilEmptyPage(t *testing.T) if err != nil { t.Fatal(err) } - list, cursor, err := client.ReadAllCurrentTasks(context.Background()) + list, cursor, err := client.ReadAllTasks(context.Background()) if err != nil { t.Fatal(err) } @@ -44,9 +44,9 @@ func TestReadAllCurrentTasksContinuesAfterShortPageUntilEmptyPage(t *testing.T) } } -func TestReadCurrentTasksFailsOnUnavailableRefreshWithoutStaleFallback(t *testing.T) { +func TestReadTasksFailsOnUnavailableRefreshWithoutStaleFallback(t *testing.T) { var unavailable atomic.Bool - body := currentExample(t, "task-discovery-page") + body := example(t, "task-discovery-page") server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") if unavailable.Load() { @@ -61,19 +61,19 @@ func TestReadCurrentTasksFailsOnUnavailableRefreshWithoutStaleFallback(t *testin if err != nil { t.Fatal(err) } - page, err := client.ReadCurrentTasks(context.Background(), "") + page, err := client.ReadTasks(context.Background(), "") if err != nil || len(page.Tasks) != 1 { t.Fatalf("expected a fresh assigned task before outage: page=%+v err=%v", page, err) } unavailable.Store(true) - _, err = client.ReadCurrentTasks(context.Background(), page.Cursor) + _, err = client.ReadTasks(context.Background(), page.Cursor) var httpErr *HTTPError if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable || !strings.Contains(err.Error(), "task discovery") { t.Fatalf("refresh failure lost its owner/status or reused stale data: %v", err) } } -func TestReadCurrentTasksFailClosed(t *testing.T) { +func TestReadTasksFailClosed(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" for _, tc := range []struct { name string @@ -100,7 +100,7 @@ func TestReadCurrentTasksFailClosed(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := client.ReadCurrentTasks(context.Background(), tc.after); err == nil { + if _, err := client.ReadTasks(context.Background(), tc.after); err == nil { t.Fatal("unexpectedly accepted invalid discovery page") } }) diff --git a/internal/configread/current.go b/internal/configread/snapshots.go similarity index 59% rename from internal/configread/current.go rename to internal/configread/snapshots.go index ab8a3ac..da7f0d6 100644 --- a/internal/configread/current.go +++ b/internal/configread/snapshots.go @@ -11,23 +11,23 @@ import ( "git.ipao.vip/rogee/go-sip/internal/contract" ) -// CurrentSnapshot carries the four independently validated configuration +// Snapshot carries the four independently validated configuration // resources. Agent.Raw preserves explicit zero/false and immutable AI fields. -type CurrentSnapshot struct { - SIP CurrentSIP - Task CurrentTask - Quota CurrentQuota - Providers map[string]CurrentProvider +type Snapshot struct { + SIP SIP + Task Task + Quota Quota + Providers map[string]Provider } -type CurrentSIP struct { +type SIP struct { Resource string `json:"resource"` DispatcherID string `json:"dispatcher_id"` Revision int64 `json:"revision"` Trunks json.RawMessage `json:"trunks"` } -type CurrentProvider struct { +type Provider struct { ProviderRef string `json:"provider_ref"` Role string `json:"role"` Enabled bool `json:"enabled"` @@ -36,13 +36,13 @@ type CurrentProvider struct { Credential string `json:"credential"` } -type currentProviderResponse struct { - Resource string `json:"resource"` - DispatcherID string `json:"dispatcher_id"` - Providers []CurrentProvider `json:"providers"` +type providerResponse struct { + Resource string `json:"resource"` + DispatcherID string `json:"dispatcher_id"` + Providers []Provider `json:"providers"` } -type CurrentAgent struct { +type Agent struct { Mode string `json:"mode"` ASR struct { ProviderRef string `json:"provider_ref"` @@ -56,18 +56,18 @@ type CurrentAgent struct { Raw json.RawMessage `json:"-"` } -func (a *CurrentAgent) UnmarshalJSON(raw []byte) error { - type plain CurrentAgent +func (a *Agent) UnmarshalJSON(raw []byte) error { + type plain Agent var parsed plain if err := json.Unmarshal(raw, &parsed); err != nil { return err } - *a = CurrentAgent(parsed) + *a = Agent(parsed) a.Raw = append([]byte(nil), raw...) return nil } -type CurrentTask struct { +type Task struct { Resource string `json:"resource"` DispatcherID string `json:"dispatcher_id"` TenantID int64 `json:"tenant_id"` @@ -82,22 +82,22 @@ type CurrentTask struct { CallerProfileID string `json:"caller_profile_id"` AllowedTrunkIDs []string `json:"allowed_trunk_ids"` Schedule json.RawMessage `json:"schedule"` - Agent CurrentAgent `json:"agent"` + Agent Agent `json:"agent"` Raw json.RawMessage `json:"-"` } -func (t *CurrentTask) UnmarshalJSON(raw []byte) error { - type plain CurrentTask +func (t *Task) UnmarshalJSON(raw []byte) error { + type plain Task var parsed plain if err := json.Unmarshal(raw, &parsed); err != nil { return err } - *t = CurrentTask(parsed) + *t = Task(parsed) t.Raw = append([]byte(nil), raw...) return nil } -type CurrentQuota struct { +type Quota struct { Resource string `json:"resource"` DispatcherID string `json:"dispatcher_id"` TenantID int64 `json:"tenant_id"` @@ -105,49 +105,49 @@ type CurrentQuota struct { MaxConcurrentCalls int64 `json:"max_concurrent_calls"` } -// ReadCurrentSIP is always called at Dispatcher startup, even with no tasks. +// ReadSIP is always called at Dispatcher startup, even with no tasks. // Revision is not a claim that Agent/Asterisk actually loaded this snapshot; // the caller must verify applied state before opening admission. -func (c *Client) ReadCurrentSIP(ctx context.Context) (CurrentSIP, error) { - var sip CurrentSIP - if err := c.readCurrent(ctx, configReadPath+"/sip", "sip_config", &sip); err != nil { - return CurrentSIP{}, fmt.Errorf("read SIP configuration: %w", err) +func (c *Client) ReadSIP(ctx context.Context) (SIP, error) { + var sip SIP + if err := c.readResource(ctx, configReadPath+"/sip", "sip_config", &sip); err != nil { + return SIP{}, fmt.Errorf("read SIP configuration: %w", err) } if sip.DispatcherID != c.dispatcherID { - return CurrentSIP{}, errors.New("SIP configuration dispatcher owner mismatch") + return SIP{}, errors.New("SIP configuration dispatcher owner mismatch") } return sip, nil } -// ReadCurrentTask reads only the agreed read-only HTTP resources. It never +// ReadTask reads only the agreed read-only HTTP resources. It never // substitutes stale data or the historical MQ configuration contract. -func (c *Client) ReadCurrentTask(ctx context.Context, taskID string, tenantID int64) (CurrentSnapshot, error) { +func (c *Client) ReadTask(ctx context.Context, taskID string, tenantID int64) (Snapshot, error) { if taskID == "" || tenantID <= 0 { - return CurrentSnapshot{}, errors.New("task ID and positive tenant ID are required") + return Snapshot{}, errors.New("task ID and positive tenant ID are required") } - var result CurrentSnapshot + var result Snapshot var err error - result.SIP, err = c.ReadCurrentSIP(ctx) + result.SIP, err = c.ReadSIP(ctx) if err != nil { - return CurrentSnapshot{}, err + return Snapshot{}, err } - var providers currentProviderResponse - if err := c.readCurrent(ctx, configReadPath+"/ai-providers", "ai_providers", &providers); err != nil { - return CurrentSnapshot{}, fmt.Errorf("read AI providers: %w", err) + var providers providerResponse + if err := c.readResource(ctx, configReadPath+"/ai-providers", "ai_providers", &providers); err != nil { + return Snapshot{}, fmt.Errorf("read AI providers: %w", err) } - if err := c.readCurrent(ctx, configReadPath+"/task/"+url.PathEscape(taskID), "task_config", &result.Task); err != nil { - return CurrentSnapshot{}, fmt.Errorf("read task configuration: %w", err) + if err := c.readResource(ctx, configReadPath+"/task/"+url.PathEscape(taskID), "task_config", &result.Task); err != nil { + return Snapshot{}, fmt.Errorf("read task configuration: %w", err) } - if err := c.readCurrent(ctx, configReadPath+"/tenant/"+strconv.FormatInt(tenantID, 10)+"/quota", "tenant_quota", &result.Quota); err != nil { - return CurrentSnapshot{}, fmt.Errorf("read tenant quota: %w", err) + if err := c.readResource(ctx, configReadPath+"/tenant/"+strconv.FormatInt(tenantID, 10)+"/quota", "tenant_quota", &result.Quota); err != nil { + return Snapshot{}, fmt.Errorf("read tenant quota: %w", err) } if result.SIP.DispatcherID != c.dispatcherID || providers.DispatcherID != c.dispatcherID || result.Task.DispatcherID != c.dispatcherID || result.Quota.DispatcherID != c.dispatcherID || result.Task.TenantID != tenantID || result.Quota.TenantID != tenantID || result.Task.TaskID != taskID { - return CurrentSnapshot{}, errors.New("configuration dispatcher, task, or tenant owner mismatch") + return Snapshot{}, errors.New("configuration dispatcher, task, or tenant owner mismatch") } - result.Providers = make(map[string]CurrentProvider, len(providers.Providers)) + result.Providers = make(map[string]Provider, len(providers.Providers)) for _, provider := range providers.Providers { if _, exists := result.Providers[provider.ProviderRef]; exists { - return CurrentSnapshot{}, fmt.Errorf("duplicate AI provider reference %q", provider.ProviderRef) + return Snapshot{}, fmt.Errorf("duplicate AI provider reference %q", provider.ProviderRef) } result.Providers[provider.ProviderRef] = provider } @@ -161,13 +161,13 @@ func (c *Client) ReadCurrentTask(ctx context.Context, taskID string, tenantID in } provider, found := result.Providers[expected.ref] if !found || !provider.Enabled || provider.Role != expected.role { - return CurrentSnapshot{}, fmt.Errorf("AI provider %q is missing, disabled, or has an incompatible role", expected.ref) + return Snapshot{}, fmt.Errorf("AI provider %q is missing, disabled, or has an incompatible role", expected.ref) } } return result, nil } -func (c *Client) readCurrent(ctx context.Context, path, resource string, dst any) error { +func (c *Client) readResource(ctx context.Context, path, resource string, dst any) error { body, err := c.get(ctx, path) if err != nil { return err diff --git a/internal/configread/current_test.go b/internal/configread/snapshots_test.go similarity index 74% rename from internal/configread/current_test.go rename to internal/configread/snapshots_test.go index 644db42..8267b64 100644 --- a/internal/configread/current_test.go +++ b/internal/configread/snapshots_test.go @@ -13,7 +13,7 @@ import ( "testing" ) -func currentExample(t *testing.T, name string) []byte { +func example(t *testing.T, name string) []byte { t.Helper() raw, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", name+".json")) if err != nil { @@ -22,14 +22,14 @@ func currentExample(t *testing.T, name string) []byte { return raw } -func TestReadCurrentTaskSnapshot(t *testing.T) { +func TestReadTaskSnapshot(t *testing.T) { for _, mode := range []string{"asr", "full"} { t.Run(mode, func(t *testing.T) { responses := map[string][]byte{ - "/internal/v1/dispatcher/sip": currentExample(t, "config-read-sip"), - "/internal/v1/dispatcher/ai-providers": currentExample(t, "config-read-providers"), - "/internal/v1/dispatcher/task/task-" + mode: currentExample(t, "config-read-task-"+mode), - "/internal/v1/dispatcher/tenant/1001/quota": currentExample(t, "config-read-quota"), + "/internal/v1/dispatcher/sip": example(t, "config-read-sip"), + "/internal/v1/dispatcher/ai-providers": example(t, "config-read-providers"), + "/internal/v1/dispatcher/task/task-" + mode: example(t, "config-read-task-"+mode), + "/internal/v1/dispatcher/tenant/1001/quota": example(t, "config-read-quota"), } calls := 0 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -53,7 +53,7 @@ func TestReadCurrentTaskSnapshot(t *testing.T) { if err != nil { t.Fatal(err) } - snapshot, err := client.ReadCurrentTask(context.Background(), "task-"+mode, 1001) + snapshot, err := client.ReadTask(context.Background(), "task-"+mode, 1001) if err != nil { t.Fatal(err) } @@ -83,7 +83,7 @@ func TestReadCurrentTaskSnapshot(t *testing.T) { } } -func TestReadCurrentSIPWithoutTasks(t *testing.T) { +func TestReadSIPWithoutTasks(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" requests := 0 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -92,26 +92,26 @@ func TestReadCurrentSIPWithoutTasks(t *testing.T) { t.Errorf("startup SIP read requested unexpected path %q", r.URL.Path) } w.Header().Set("Content-Type", "application/json") - _, _ = w.Write(currentExample(t, "config-read-sip")) + _, _ = w.Write(example(t, "config-read-sip")) })) defer server.Close() client, err := NewClient(server.URL, id, "test-secret", server.Client()) if err != nil { t.Fatal(err) } - sip, err := client.ReadCurrentSIP(context.Background()) + sip, err := client.ReadSIP(context.Background()) if err != nil || sip.Revision != 8 || requests != 1 { t.Fatalf("startup SIP config: %+v, %v, requests=%d", sip, err, requests) } } -func TestReadCurrentTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) { +func TestReadTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) { const dispatcherID = "c046b893-8628-4589-ae50-619d049248a6" responses := map[string][]byte{ - "/internal/v1/dispatcher/sip": currentExample(t, "config-read-sip"), - "/internal/v1/dispatcher/ai-providers": currentExample(t, "config-read-providers"), - "/internal/v1/dispatcher/task/task-asr": currentExample(t, "config-read-task-asr"), - "/internal/v1/dispatcher/tenant/1001/quota": currentExample(t, "config-read-quota"), + "/internal/v1/dispatcher/sip": example(t, "config-read-sip"), + "/internal/v1/dispatcher/ai-providers": example(t, "config-read-providers"), + "/internal/v1/dispatcher/task/task-asr": example(t, "config-read-task-asr"), + "/internal/v1/dispatcher/tenant/1001/quota": example(t, "config-read-quota"), } for _, tc := range []struct{ path, owner string }{ {"/internal/v1/dispatcher/sip", "SIP"}, @@ -140,11 +140,11 @@ func TestReadCurrentTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := client.ReadCurrentTask(context.Background(), "task-asr", 1001); err != nil { + if _, err := client.ReadTask(context.Background(), "task-asr", 1001); err != nil { t.Fatalf("approved resources failed before the outage: %v", err) } unavailable.Store(true) - _, err = client.ReadCurrentTask(context.Background(), "task-asr", 1001) + _, err = client.ReadTask(context.Background(), "task-asr", 1001) var httpErr *HTTPError if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable { t.Fatalf("unavailable %s reused cached configuration or lost status: %v", tc.owner, err) @@ -156,22 +156,22 @@ func TestReadCurrentTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) { } } -func TestReadCurrentTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) { +func TestReadTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) { for _, mode := range []string{"owner", "sip-revision", "provider-ref", "provider-disabled", "provider-wrong-role"} { t.Run(mode, func(t *testing.T) { responses := map[string][]byte{ - "/internal/v1/dispatcher/sip": currentExample(t, "config-read-sip"), - "/internal/v1/dispatcher/ai-providers": currentExample(t, "config-read-providers"), - "/internal/v1/dispatcher/task/task-asr": currentExample(t, "config-read-task-asr"), - "/internal/v1/dispatcher/tenant/1001/quota": currentExample(t, "config-read-quota"), + "/internal/v1/dispatcher/sip": example(t, "config-read-sip"), + "/internal/v1/dispatcher/ai-providers": example(t, "config-read-providers"), + "/internal/v1/dispatcher/task/task-asr": example(t, "config-read-task-asr"), + "/internal/v1/dispatcher/tenant/1001/quota": example(t, "config-read-quota"), } switch mode { case "owner": responses["/internal/v1/dispatcher/task/task-asr"] = []byte(strings.Replace(string(responses["/internal/v1/dispatcher/task/task-asr"]), `"tenant_id":1001`, `"tenant_id":1002`, 1)) case "sip-revision": - responses["/internal/v1/dispatcher/sip"] = currentExample(t, "invalid/config-read-sip-missing-revision") + responses["/internal/v1/dispatcher/sip"] = example(t, "invalid/config-read-sip-missing-revision") case "provider-ref": - responses["/internal/v1/dispatcher/ai-providers"] = currentExample(t, "invalid/config-read-provider-ref") + responses["/internal/v1/dispatcher/ai-providers"] = example(t, "invalid/config-read-provider-ref") case "provider-disabled": responses["/internal/v1/dispatcher/ai-providers"] = []byte(strings.Replace(string(responses["/internal/v1/dispatcher/ai-providers"]), `"enabled":true`, `"enabled":false`, 1)) case "provider-wrong-role": @@ -186,7 +186,7 @@ func TestReadCurrentTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := client.ReadCurrentTask(context.Background(), "task-asr", 1001); err == nil || strings.Contains(err.Error(), "example-only-not-a-real-secret") { + if _, err := client.ReadTask(context.Background(), "task-asr", 1001); err == nil || strings.Contains(err.Error(), "example-only-not-a-real-secret") { t.Fatalf("expected redacted fail-closed error, got %v", err) } }) diff --git a/internal/dispatcher/approved_originator.go b/internal/dispatcher/approved_originator.go index 2ece692..9865388 100644 --- a/internal/dispatcher/approved_originator.go +++ b/internal/dispatcher/approved_originator.go @@ -70,7 +70,7 @@ func (o *ApprovedOriginator) LoadedTrunks(ctx context.Context) (map[string]int64 // VerifySIP requires an exact match between the approved entire SIP partition // and the Agent's actual loaded set; a single matching trunk is insufficient. -func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.CurrentSIP) error { +func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.SIP) error { if sip.DispatcherID != o.DispatcherID || sip.Revision <= 0 { return errors.New("approved SIP snapshot does not belong to this Dispatcher") } diff --git a/internal/dispatcher/approved_originator_test.go b/internal/dispatcher/approved_originator_test.go index 4207e3d..f3a02af 100644 --- a/internal/dispatcher/approved_originator_test.go +++ b/internal/dispatcher/approved_originator_test.go @@ -88,7 +88,7 @@ func TestApprovedOriginatorSendsExactImmutableAgentInstruction(t *testing.T) { if !bytes.Equal(req.TaskConfigJson, spec.Snapshot.Task.Raw) { t.Fatal("Agent did not receive the exact immutable task JSON") } - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_ai.go b/internal/dispatcher/current_ai.go index 26b8e09..3ad747c 100644 --- a/internal/dispatcher/current_ai.go +++ b/internal/dispatcher/current_ai.go @@ -10,7 +10,7 @@ import ( // validateCurrentAISnapshot is the shared admission gate for every source of // approved task configuration. SDK-inexpressible settings must fail before a // call reserves capacity or reaches the Agent; no defaults replace them. -func validateCurrentAISnapshot(snapshot configread.CurrentSnapshot) error { +func validateCurrentAISnapshot(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/current_ai_test.go index d21975a..3759f9e 100644 --- a/internal/dispatcher/current_ai_test.go +++ b/internal/dispatcher/current_ai_test.go @@ -63,7 +63,7 @@ func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *test drained := false b := CurrentBootstrap{ DispatcherID: id, Client: client, Store: db, - VerifySIP: func(context.Context, configread.CurrentSIP) error { return nil }, + VerifySIP: func(context.Context, configread.SIP) error { return nil }, DrainControls: func(context.Context) error { drained = true; return nil }, } err = b.Run(context.Background()) @@ -89,7 +89,7 @@ func TestCurrentExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing. t.Fatal(err) } id := snapshot.Task.DispatcherID - if err := db.ApplyDiscoverySnapshot(id, []configread.CurrentDiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { + if err := db.ApplyDiscoverySnapshot(id, []configread.DiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { t.Fatal(err) } if err := db.SaveSnapshot(snapshot); err != nil { diff --git a/internal/dispatcher/current_config.go b/internal/dispatcher/current_config.go index ca4b147..0ca5868 100644 --- a/internal/dispatcher/current_config.go +++ b/internal/dispatcher/current_config.go @@ -18,10 +18,10 @@ type CurrentBootstrap struct { Client *configread.Client Store *store.Store DispatcherID string - VerifySIP func(context.Context, configread.CurrentSIP) error + VerifySIP func(context.Context, configread.SIP) error DrainControls func(context.Context) error Cursor *string // memory-only; restart always starts from a full snapshot - SIP *configread.CurrentSIP // the exact approved snapshot verified on startup + SIP *configread.SIP // the exact approved snapshot verified on startup } func (b CurrentBootstrap) Run(ctx context.Context) error { @@ -32,12 +32,12 @@ func (b CurrentBootstrap) Run(ctx context.Context) error { *b.Cursor = "" } if b.SIP != nil { - *b.SIP = configread.CurrentSIP{} + *b.SIP = configread.SIP{} } if err := b.Store.CloseAdmission(b.DispatcherID); err != nil { return fmt.Errorf("close task admission before bootstrap: %w", err) } - sip, err := b.Client.ReadCurrentSIP(ctx) + sip, err := b.Client.ReadSIP(ctx) if err != nil { return fmt.Errorf("read approved SIP snapshot: %w", err) } @@ -47,7 +47,7 @@ func (b CurrentBootstrap) Run(ctx context.Context) error { if err := b.VerifySIP(ctx, sip); err != nil { return fmt.Errorf("verify SIP revision %d loaded by Agent/Asterisk: %w", sip.Revision, err) } - tasks, cursor, err := b.Client.ReadAllCurrentTasks(ctx) + tasks, cursor, err := b.Client.ReadAllTasks(ctx) if err != nil { return fmt.Errorf("retrieve complete assigned task list: %w", err) } @@ -58,7 +58,7 @@ func (b CurrentBootstrap) Run(ctx context.Context) error { if task.Status == "stopped" { continue } - snapshot, err := b.Client.ReadCurrentTask(ctx, task.TaskID, task.TenantID) + snapshot, err := b.Client.ReadTask(ctx, task.TaskID, task.TenantID) if err != nil { return fmt.Errorf("read assigned task %q: %w", task.TaskID, err) } diff --git a/internal/dispatcher/current_config_test.go b/internal/dispatcher/current_config_test.go index 54152fc..a80e792 100644 --- a/internal/dispatcher/current_config_test.go +++ b/internal/dispatcher/current_config_test.go @@ -63,10 +63,10 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { defer db.Close() verifierCalled, drained := false, false var cursor string - var approvedSIP configread.CurrentSIP + var approvedSIP configread.SIP bootstrap := CurrentBootstrap{ Client: client, Store: db, DispatcherID: id, Cursor: &cursor, SIP: &approvedSIP, - VerifySIP: func(_ context.Context, sip configread.CurrentSIP) error { + VerifySIP: func(_ context.Context, sip configread.SIP) error { verifierCalled = true if sip.Revision != 8 { t.Fatalf("verified wrong SIP revision %d", sip.Revision) @@ -141,7 +141,7 @@ func TestCurrentBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { defer db.Close() b := CurrentBootstrap{ Client: client, Store: db, DispatcherID: id, - VerifySIP: func(context.Context, configread.CurrentSIP) error { + VerifySIP: func(context.Context, configread.SIP) error { if failure == "verify" { return errors.New("applied SIP revision unavailable") } diff --git a/internal/dispatcher/current_control.go b/internal/dispatcher/current_control.go index a00a96a..c43890c 100644 --- a/internal/dispatcher/current_control.go +++ b/internal/dispatcher/current_control.go @@ -32,7 +32,7 @@ type CurrentControlController struct { Store *store.Store Client *configread.Client Agent CurrentControlAgent - VerifySIP func(context.Context, configread.CurrentSIP) error + VerifySIP func(context.Context, configread.SIP) error Now func() time.Time } @@ -83,7 +83,7 @@ func (c *CurrentControlController) ProcessControl(ctx context.Context, body []by // A stale HTTP running status cannot override persisted pause/stop. // Only a fresh approved task plus applied SIP snapshot may authorize // the transition from a durable pause back to admission. - snapshot, err := c.Client.ReadCurrentTask(ctx, event.Payload.TaskID, event.TenantID) + snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID) if err != nil { return fmt.Errorf("fresh resume task configuration: %w", err) } diff --git a/internal/dispatcher/current_control_test.go b/internal/dispatcher/current_control_test.go index bd99ef5..10844a2 100644 --- a/internal/dispatcher/current_control_test.go +++ b/internal/dispatcher/current_control_test.go @@ -68,7 +68,7 @@ func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *current t.Fatal(err) } t.Cleanup(func() { _ = s.Close() }) - if err := s.ApplyDiscoverySnapshot(id, []configread.CurrentDiscoveredTask{{TaskID: "task-asr", TenantID: 1001, TaskRevision: 1, Status: "running"}}); err != nil { + if err := s.ApplyDiscoverySnapshot(id, []configread.DiscoveredTask{{TaskID: "task-asr", TenantID: 1001, TaskRevision: 1, Status: "running"}}); err != nil { t.Fatal(err) } if err := s.SaveSnapshot(snapshot); err != nil { @@ -78,7 +78,7 @@ func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *current t.Fatal(err) } agent := ¤tFakeControlAgent{} - controller := &CurrentControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, VerifySIP: func(context.Context, configread.CurrentSIP) error { return nil }, Now: func() time.Time { return currentMonday(9, 30) }} + 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) }} return controller, agent, s } @@ -185,7 +185,7 @@ func TestCurrentResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) { t.Fatal(err) } approved := controller.VerifySIP - controller.VerifySIP = func(context.Context, configread.CurrentSIP) error { + 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 { diff --git a/internal/dispatcher/current_discovery.go b/internal/dispatcher/current_discovery.go index 870733b..bab3290 100644 --- a/internal/dispatcher/current_discovery.go +++ b/internal/dispatcher/current_discovery.go @@ -17,8 +17,8 @@ type CurrentDiscoveryFollower struct { DispatcherID string Client *configread.Client Store *store.Store - ApprovedSIP configread.CurrentSIP - VerifySIP func(context.Context, configread.CurrentSIP) error + ApprovedSIP configread.SIP + VerifySIP func(context.Context, configread.SIP) error Cursor string } @@ -31,7 +31,7 @@ func (f *CurrentDiscoveryFollower) Poll(ctx context.Context) error { return f.fail(fmt.Errorf("encode approved SIP snapshot: %w", err)) } for { - page, err := f.Client.ReadCurrentTasks(ctx, f.Cursor) + page, err := f.Client.ReadTasks(ctx, f.Cursor) if err != nil { return f.fail(fmt.Errorf("read task discovery after cursor: %w", err)) } @@ -46,7 +46,7 @@ func (f *CurrentDiscoveryFollower) Poll(ctx context.Context) error { if task.Status == "stopped" { continue } - snapshot, err := f.Client.ReadCurrentTask(ctx, task.TaskID, task.TenantID) + snapshot, err := f.Client.ReadTask(ctx, task.TaskID, task.TenantID) if err != nil { return f.fail(fmt.Errorf("load discovered task %q: %w", task.TaskID, err)) } diff --git a/internal/dispatcher/current_discovery_test.go b/internal/dispatcher/current_discovery_test.go index 666a387..c3168ab 100644 --- a/internal/dispatcher/current_discovery_test.go +++ b/internal/dispatcher/current_discovery_test.go @@ -56,7 +56,7 @@ 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.CurrentSIP) error { return nil }, Cursor: "initial"} + follower := &CurrentDiscoveryFollower{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 } } diff --git a/internal/dispatcher/current_execute.go b/internal/dispatcher/current_execute.go index ab04032..d364155 100644 --- a/internal/dispatcher/current_execute.go +++ b/internal/dispatcher/current_execute.go @@ -27,7 +27,7 @@ type CurrentCallSpec struct { RingTimeoutMS int64 MaxCallDurationMS int64 Deadline time.Time - Snapshot configread.CurrentSnapshot + Snapshot configread.Snapshot } type CurrentOriginator interface { diff --git a/internal/dispatcher/current_execute_test.go b/internal/dispatcher/current_execute_test.go index d007b86..7427893 100644 --- a/internal/dispatcher/current_execute_test.go +++ b/internal/dispatcher/current_execute_test.go @@ -52,7 +52,7 @@ func newCurrentExecuteFixture(t *testing.T) (*CurrentExecuteController, *current } t.Cleanup(func() { _ = s.Close() }) snapshot := currentPolicySnapshot(t) - if err := s.ApplyDiscoverySnapshot(snapshot.Task.DispatcherID, []configread.CurrentDiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { + 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) } if err := s.SaveSnapshot(snapshot); err != nil { diff --git a/internal/dispatcher/current_policy.go b/internal/dispatcher/current_policy.go index 771a3f8..8e74696 100644 --- a/internal/dispatcher/current_policy.go +++ b/internal/dispatcher/current_policy.go @@ -49,7 +49,7 @@ var currentCalleeWhitelist = map[string]struct{}{"15003164745": {}, "15830461047 // SelectCurrentTrunk 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.CurrentSnapshot, callee string, at time.Time, trunkOccupancy, loadedRevisions map[string]int64) (CurrentSelectedTrunk, error) { +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 } diff --git a/internal/dispatcher/current_policy_test.go b/internal/dispatcher/current_policy_test.go index 240ecea..c5937e4 100644 --- a/internal/dispatcher/current_policy_test.go +++ b/internal/dispatcher/current_policy_test.go @@ -10,9 +10,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func currentPolicySnapshot(t *testing.T) configread.CurrentSnapshot { +func currentPolicySnapshot(t *testing.T) configread.Snapshot { t.Helper() - var snapshot configread.CurrentSnapshot + var snapshot configread.Snapshot if err := json.Unmarshal(currentConfigExample(t, "config-read-task-asr"), &snapshot.Task); err != nil { t.Fatal(err) } @@ -27,12 +27,12 @@ func currentPolicySnapshot(t *testing.T) configread.CurrentSnapshot { t.Fatal(err) } var providerList struct { - Providers []configread.CurrentProvider `json:"providers"` + Providers []configread.Provider `json:"providers"` } if err := json.Unmarshal(currentConfigExample(t, "config-read-providers"), &providerList); err != nil { t.Fatal(err) } - snapshot.Providers = make(map[string]configread.CurrentProvider, len(providerList.Providers)) + snapshot.Providers = make(map[string]configread.Provider, len(providerList.Providers)) for _, provider := range providerList.Providers { snapshot.Providers[provider.ProviderRef] = provider } diff --git a/internal/dispatcher/current_runtime.go b/internal/dispatcher/current_runtime.go index 054d4d9..b9a2caf 100644 --- a/internal/dispatcher/current_runtime.go +++ b/internal/dispatcher/current_runtime.go @@ -67,7 +67,7 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { }() var cursor string - var sip configread.CurrentSIP + var sip configread.SIP r.Bootstrap.Cursor = &cursor r.Bootstrap.SIP = &sip r.Bootstrap.DrainControls = func(ctx context.Context) error { diff --git a/internal/dispatcher/current_runtime_integration_test.go b/internal/dispatcher/current_runtime_integration_test.go index 5ed7c94..0e2bf70 100644 --- a/internal/dispatcher/current_runtime_integration_test.go +++ b/internal/dispatcher/current_runtime_integration_test.go @@ -145,7 +145,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T agent := ¤tRuntimeMockAgent{calls: make(chan CurrentCallSpec, 3), controls: make(chan CurrentControlSpec, 3)} var windowAllowed atomic.Bool windowAllowed.Store(true) - verify := func(_ context.Context, sip configread.CurrentSIP) error { + verify := func(_ context.Context, sip configread.SIP) error { if sip.Revision != 8 || sip.DispatcherID != id { return fmt.Errorf("unloaded SIP revision") } diff --git a/internal/dispatcher/current_sip_conflict_test.go b/internal/dispatcher/current_sip_conflict_test.go index 2bd6a40..ebb0838 100644 --- a/internal/dispatcher/current_sip_conflict_test.go +++ b/internal/dispatcher/current_sip_conflict_test.go @@ -53,7 +53,7 @@ func TestCurrentBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) } defer db.Close() bootstrap := CurrentBootstrap{Client: client, Store: db, DispatcherID: id, - VerifySIP: func(context.Context, configread.CurrentSIP) error { return nil }, + VerifySIP: func(context.Context, configread.SIP) error { return nil }, DrainControls: func(context.Context) error { return nil }, } if err := bootstrap.Run(context.Background()); err == nil { diff --git a/internal/dispatcher/current_sip_reload.go b/internal/dispatcher/current_sip_reload.go index 6454d9c..8ab04f0 100644 --- a/internal/dispatcher/current_sip_reload.go +++ b/internal/dispatcher/current_sip_reload.go @@ -15,7 +15,7 @@ func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscov 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") } - current, err := r.Bootstrap.Client.ReadCurrentSIP(ctx) + current, err := r.Bootstrap.Client.ReadSIP(ctx) if err != nil { return r.closeSIPAdmission(fmt.Errorf("read approved SIP full snapshot: %w", err)) } @@ -68,7 +68,7 @@ func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscov // A full discovery snapshot is required after an approved SIP revision // change; a delta page cannot prove that all assigned tasks use the same // loaded revision. MQ control consumption remains active during this read. - tasks, cursor, err := r.Bootstrap.Client.ReadAllCurrentTasks(ctx) + tasks, cursor, err := r.Bootstrap.Client.ReadAllTasks(ctx) if err != nil { return r.closeSIPAdmission(fmt.Errorf("reload complete assigned task list for SIP: %w", err)) } @@ -79,7 +79,7 @@ func (r *CurrentRuntime) refreshSIP(ctx context.Context, follower *CurrentDiscov if task.Status == "stopped" { continue } - snapshot, err := r.Bootstrap.Client.ReadCurrentTask(ctx, task.TaskID, task.TenantID) + snapshot, err := r.Bootstrap.Client.ReadTask(ctx, task.TaskID, task.TenantID) if err != nil { return r.closeSIPAdmission(fmt.Errorf("reload task %q after SIP revision: %w", task.TaskID, err)) } diff --git a/internal/dispatcher/current_sip_runtime_test.go b/internal/dispatcher/current_sip_runtime_test.go index 403345e..a4be074 100644 --- a/internal/dispatcher/current_sip_runtime_test.go +++ b/internal/dispatcher/current_sip_runtime_test.go @@ -54,7 +54,7 @@ func TestCurrentSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *t t.Fatal(err) } var loaded atomic.Bool - verify := func(_ context.Context, sip configread.CurrentSIP) error { + verify := func(_ context.Context, sip configread.SIP) error { if sip.Revision != 9 || !loaded.Load() { return errors.New("Agent and Asterisk have not applied SIP revision 9") } diff --git a/internal/rpc/approved_execution.go b/internal/rpc/approved_execution.go index ff99042..e1323f6 100644 --- a/internal/rpc/approved_execution.go +++ b/internal/rpc/approved_execution.go @@ -101,14 +101,14 @@ func (s *Server) ExecuteApproved(ctx context.Context, req *agentpb.ExecuteApprov if err != nil || hash != req.BindingSha256 { return nil, status.Error(codes.FailedPrecondition, "bound execution snapshot does not match") } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(req.TaskConfigJson, &task); err != nil { return nil, status.Error(codes.InvalidArgument, "approved task JSON is invalid") } if task.DispatcherID != dispatcherID || task.TenantID != req.TenantId || task.TaskID != req.TaskId { return nil, status.Error(codes.FailedPrecondition, "approved task identity does not match") } - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { return nil, status.Error(codes.InvalidArgument, "approved provider JSON is invalid") } diff --git a/internal/rpc/approved_execution_test.go b/internal/rpc/approved_execution_test.go index 2ef48d5..7c05abc 100644 --- a/internal/rpc/approved_execution_test.go +++ b/internal/rpc/approved_execution_test.go @@ -22,7 +22,7 @@ func approvedTestRequest(t *testing.T, now time.Time) *agentpb.ExecuteApprovedRe if err != nil { t.Fatal(err) } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(taskJSON, &task); err != nil { t.Fatal(err) } @@ -31,12 +31,12 @@ func approvedTestRequest(t *testing.T, now time.Time) *agentpb.ExecuteApprovedRe t.Fatal(err) } var payload struct { - Providers []configread.CurrentProvider `json:"providers"` + Providers []configread.Provider `json:"providers"` } if err := json.Unmarshal(rawProviders, &payload); err != nil { t.Fatal(err) } - providers := make(map[string]configread.CurrentProvider) + providers := make(map[string]configread.Provider) for _, p := range payload.Providers { providers[p.ProviderRef] = p } diff --git a/internal/rpc/approved_full_ai_integration_test.go b/internal/rpc/approved_full_ai_integration_test.go index c380566..5085c91 100644 --- a/internal/rpc/approved_full_ai_integration_test.go +++ b/internal/rpc/approved_full_ai_integration_test.go @@ -32,7 +32,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T t.Fatal(err) } req.TaskConfigJson = fullJSON - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(fullJSON, &task); err != nil { t.Fatal(err) } @@ -60,7 +60,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T } })) defer mockSDKs.Close() - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { t.Fatal(err) } @@ -123,7 +123,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if err != nil { t.Fatal(err) } - var sip configread.CurrentSIP + var sip configread.SIP if err := json.Unmarshal(sipJSON, &sip); err != nil { t.Fatal(err) } @@ -138,7 +138,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T TrunkID: req.SelectedTrunkId, Callee: req.Callee, DialedCallee: req.DialedCallee, CallerID: req.CallerId, RingTimeoutMS: req.RingTimeoutMs, MaxCallDurationMS: req.MaxCallDurationMs, Deadline: time.UnixMilli(req.DialBeforeUnixMs), - Snapshot: configread.CurrentSnapshot{SIP: sip, Task: task, Providers: providers}, + Snapshot: configread.Snapshot{SIP: sip, Task: task, Providers: providers}, } if err := orig.Originate(context.Background(), spec); err != nil || calls != 1 || hangups != 1 { t.Fatalf("full-AI isolated Agent call failed: err=%v calls=%d hangups=%d", err, calls, hangups) diff --git a/internal/rpc/approved_integration_test.go b/internal/rpc/approved_integration_test.go index 4a8cadd..3b78398 100644 --- a/internal/rpc/approved_integration_test.go +++ b/internal/rpc/approved_integration_test.go @@ -56,19 +56,19 @@ func TestApprovedDispatcherToAgentUnaryMockRetainsSnapshotAndOneShotCall(t *test if err != nil { t.Fatal(err) } - var sip configread.CurrentSIP + var sip configread.SIP if err := json.Unmarshal(sipJSON, &sip); err != nil { t.Fatal(err) } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(req.TaskConfigJson, &task); err != nil { t.Fatal(err) } - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { t.Fatal(err) } - snapshot := configread.CurrentSnapshot{SIP: sip, Task: task, Providers: providers} + snapshot := configread.Snapshot{SIP: sip, Task: task, Providers: providers} if err := orig.VerifySIP(context.Background(), sip); err != nil { t.Fatal(err) } diff --git a/internal/rpc/approved_runner_test.go b/internal/rpc/approved_runner_test.go index 839f4b8..d77361f 100644 --- a/internal/rpc/approved_runner_test.go +++ b/internal/rpc/approved_runner_test.go @@ -41,11 +41,11 @@ func (*blockedApprovedMedia) Stats() media.RTPStats { return media.RTPStats{} } func TestApprovedCallRunnerASROnlyHonorsSignedTimeout(t *testing.T) { req := approvedTestRequest(t, time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC)) - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(req.TaskConfigJson, &task); err != nil { t.Fatal(err) } - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { t.Fatal(err) } @@ -70,12 +70,12 @@ func TestApprovedCallRunnerSynthesizesOpeningBeforeMediaCapture(t *testing.T) { if err != nil { t.Fatal(err) } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(raw, &task); err != nil { t.Fatal(err) } req := approvedTestRequest(t, time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC)) - var providers map[string]configread.CurrentProvider + var providers map[string]configread.Provider if err := json.Unmarshal(req.ProvidersJson, &providers); err != nil { t.Fatal(err) } diff --git a/internal/rpc/recording_server_flow_test.go b/internal/rpc/recording_server_flow_test.go index 6ec4f58..b6bc5d3 100644 --- a/internal/rpc/recording_server_flow_test.go +++ b/internal/rpc/recording_server_flow_test.go @@ -26,7 +26,7 @@ import ( const recordingDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" -func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context.Context, *agentpb.RequestRecordingUploadRequest, configread.CurrentSnapshot) { +func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context.Context, *agentpb.RequestRecordingUploadRequest, configread.Snapshot) { t.Helper() database, err := store.Open(filepath.Join(t.TempDir(), "recordings.db")) if err != nil { @@ -43,20 +43,20 @@ func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context. t.Fatal(err) } } - var snapshot configread.CurrentSnapshot + var snapshot configread.Snapshot readExample("config-read-task-asr", &snapshot.Task) readExample("config-read-sip", &snapshot.SIP) readExample("config-read-quota", &snapshot.Quota) var providers struct { - Providers []configread.CurrentProvider `json:"providers"` + Providers []configread.Provider `json:"providers"` } readExample("config-read-providers", &providers) - snapshot.Providers = make(map[string]configread.CurrentProvider) + snapshot.Providers = make(map[string]configread.Provider) for _, provider := range providers.Providers { snapshot.Providers[provider.ProviderRef] = provider } snapshot.SIP.Trunks = []byte(strings.Replace(string(snapshot.SIP.Trunks), `"max_concurrent_calls":null`, `"max_concurrent_calls":2`, 1)) - if err := database.ApplyDiscoverySnapshot(recordingDispatcherID, []configread.CurrentDiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { + if err := database.ApplyDiscoverySnapshot(recordingDispatcherID, []configread.DiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { t.Fatal(err) } if err := database.SaveSnapshot(snapshot); err != nil { @@ -104,7 +104,7 @@ func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context. return server, database, ctx, request, snapshot } -func recordingResultPayload(t *testing.T, task configread.CurrentSnapshot, grant *agentpb.UploadGrant) []byte { +func recordingResultPayload(t *testing.T, task configread.Snapshot, grant *agentpb.UploadGrant) []byte { t.Helper() result := map[string]any{ "task_id": "task-asr", "caller_profile_id": task.Task.CallerProfileID, diff --git a/internal/store/calls.go b/internal/store/calls.go index 1953f89..7cf3dc7 100644 --- a/internal/store/calls.go +++ b/internal/store/calls.go @@ -144,9 +144,9 @@ func (s *Store) ReserveExecute(dispatcherID, eventID string, selected CallReserv return fmt.Errorf("%w: SIP revision changed", ErrNotReady) } var snapshot struct { - Task configread.CurrentTask `json:"task"` - SIP configread.CurrentSIP `json:"sip"` - Quota configread.CurrentQuota `json:"quota"` + Task configread.Task `json:"task"` + SIP configread.SIP `json:"sip"` + Quota configread.Quota `json:"quota"` } if err := json.Unmarshal(body, &snapshot); err != nil { return fmt.Errorf("decode admission snapshot: %w", err) diff --git a/internal/store/calls_test.go b/internal/store/calls_test.go index 1242850..3760ba9 100644 --- a/internal/store/calls_test.go +++ b/internal/store/calls_test.go @@ -19,7 +19,7 @@ func preparedCurrentCallStore(t *testing.T) *Store { t.Cleanup(func() { _ = s.Close() }) task := currentStoreSnapshot(t) task.SIP.Trunks = []byte(strings.Replace(string(task.SIP.Trunks), `"max_concurrent_calls":null`, `"max_concurrent_calls":2`, 1)) - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{{TaskID: task.Task.TaskID, TenantID: task.Task.TenantID, TaskRevision: task.Task.TaskRevision, Status: "running"}}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{{TaskID: task.Task.TaskID, TenantID: task.Task.TenantID, TaskRevision: task.Task.TaskRevision, Status: "running"}}); err != nil { t.Fatal(err) } if err := s.SaveSnapshot(task); err != nil { diff --git a/internal/store/discovery_test.go b/internal/store/discovery_test.go index 8833bd0..11b158d 100644 --- a/internal/store/discovery_test.go +++ b/internal/store/discovery_test.go @@ -7,8 +7,8 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func currentDiscoveredTask(status string, revision int64) configread.CurrentDiscoveredTask { - return configread.CurrentDiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: revision, Status: status} +func currentDiscoveredTask(status string, revision int64) configread.DiscoveredTask { + return configread.DiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: revision, Status: status} } func TestHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testing.T) { @@ -18,7 +18,7 @@ func TestHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testing.T) { } defer s.Close() base := currentStoreSnapshot(t) - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { t.Fatal(err) } if err := s.SaveSnapshot(base); err != nil { @@ -30,13 +30,13 @@ func TestHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testing.T) { if _, _, err := s.RecordExecute(currentCall("not-admitted-before-http-stop")); err != nil { t.Fatal(err) } - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("stopped", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("stopped", 1)}); err != nil { t.Fatal(err) } if pending, err := s.ListPendingExecute(currentDispatcherID); err != nil || len(pending) != 0 { t.Fatalf("HTTP stop retained pending command: %+v %v", pending, err) } - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { t.Fatal(err) } if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil { @@ -57,16 +57,16 @@ func TestHTTPPausedRequiresExplicitMQResume(t *testing.T) { } defer s.Close() base := currentStoreSnapshot(t) - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { t.Fatal(err) } if err := s.SaveSnapshot(base); err != nil { t.Fatal(err) } - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("paused", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("paused", 1)}); err != nil { t.Fatal(err) } - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { t.Fatal(err) } if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil { @@ -92,13 +92,13 @@ func TestDiscoveryPageCommitsAtomicallyAndPreservesControl(t *testing.T) { t.Fatal(err) } defer s.Close() - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{currentDiscoveredTask("running", 1)}); err != nil { t.Fatal(err) } if err := s.ApplyControl(currentDispatcherID, 1001, "task-asr", "pause"); err != nil { t.Fatal(err) } - tasks := []configread.CurrentDiscoveredTask{currentDiscoveredTask("running", 2), {TaskID: "task-second", TenantID: 0, TaskRevision: 1, Status: "running"}} + tasks := []configread.DiscoveredTask{currentDiscoveredTask("running", 2), {TaskID: "task-second", TenantID: 0, TaskRevision: 1, Status: "running"}} if err := s.ApplyDiscoveryPage(currentDispatcherID, tasks); err == nil { t.Fatal("invalid page was partly committed") } diff --git a/internal/store/global_revision_test.go b/internal/store/global_revision_test.go index 2d08a84..fb243a4 100644 --- a/internal/store/global_revision_test.go +++ b/internal/store/global_revision_test.go @@ -71,7 +71,7 @@ func TestSharedGlobalRevisionConflictAcrossTasks(t *testing.T) { t.Fatal(err) } second := first - var task configread.CurrentTask + var task configread.Task raw := []byte(strings.Replace(string(first.Task.Raw), `"task_id":"task-asr"`, `"task_id":"task-second"`, 1)) if string(raw) == string(first.Task.Raw) { t.Fatal("test task clone fixture failed") diff --git a/internal/store/result.go b/internal/store/result.go index 7ec68e6..7c0ed38 100644 --- a/internal/store/result.go +++ b/internal/store/result.go @@ -80,7 +80,7 @@ func (s *Store) recordCallResult(dispatcherID, sourceEventID string, payload []b return OutboxEvent{}, false, ErrEndUnconfirmed } var snapshot struct { - Task configread.CurrentTask `json:"task"` + Task configread.Task `json:"task"` } if err := json.Unmarshal(snapshotJSON, &snapshot); err != nil { return OutboxEvent{}, false, fmt.Errorf("decode frozen call task: %w", err) diff --git a/internal/store/store.go b/internal/store/store.go index ea5f432..91e272f 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -223,7 +223,7 @@ func (s *Store) CloseAdmission(dispatcherID string) error { // ApplyDiscoverySnapshot commits the entire cold-start list in one transaction. // It closes admission until the caller has drained the control queue. -func (s *Store) ApplyDiscoverySnapshot(dispatcherID string, tasks []configread.CurrentDiscoveredTask) (err error) { +func (s *Store) ApplyDiscoverySnapshot(dispatcherID string, tasks []configread.DiscoveredTask) (err error) { if dispatcherID == "" { return errors.New("dispatcher ID is required") } @@ -247,7 +247,7 @@ func (s *Store) ApplyDiscoverySnapshot(dispatcherID string, tasks []configread.C // ApplyDiscoveryPage commits one delta page before its in-memory cursor may // advance. It never treats tasks absent from a delta page as retired. -func (s *Store) ApplyDiscoveryPage(dispatcherID string, tasks []configread.CurrentDiscoveredTask) error { +func (s *Store) ApplyDiscoveryPage(dispatcherID string, tasks []configread.DiscoveredTask) error { if dispatcherID == "" || len(tasks) == 0 { return errors.New("discovery delta requires a Dispatcher and nonempty page") } @@ -265,7 +265,7 @@ func (s *Store) ApplyDiscoveryPage(dispatcherID string, tasks []configread.Curre return nil } -func applyDiscoveredTasks(tx *sql.Tx, dispatcherID string, tasks []configread.CurrentDiscoveredTask) error { +func applyDiscoveredTasks(tx *sql.Tx, dispatcherID string, tasks []configread.DiscoveredTask) error { seen := make(map[string]struct{}, len(tasks)) for _, task := range tasks { if task.TaskID == "" || task.TenantID <= 0 || task.TaskRevision <= 0 || !validTaskStatus(task.Status) { @@ -368,12 +368,12 @@ func (s *Store) CanAdmit(dispatcherID string, tenantID int64, taskID string) (bo // SaveSnapshot refuses an immutable task revision with different content. Only // the providers actually referenced by this task are included in its digest. -func (s *Store) SaveSnapshot(snapshot configread.CurrentSnapshot) error { +func (s *Store) SaveSnapshot(snapshot configread.Snapshot) error { task := snapshot.Task if task.DispatcherID == "" || task.TenantID <= 0 || task.TaskID == "" || task.TaskRevision <= 0 || snapshot.SIP.DispatcherID != task.DispatcherID || snapshot.Quota.DispatcherID != task.DispatcherID || snapshot.Quota.TenantID != task.TenantID || snapshot.SIP.Revision <= 0 || snapshot.Quota.QuotaRevision <= 0 || !json.Valid(task.Raw) { return errors.New("invalid current task snapshot identity, revision, or body") } - providers := make(map[string]configread.CurrentProvider) + providers := make(map[string]configread.Provider) for _, ref := range []string{task.Agent.ASR.ProviderRef, task.Agent.LLM.ProviderRef, task.Agent.TTS.ProviderRef} { if ref == "" { continue @@ -392,16 +392,16 @@ func (s *Store) SaveSnapshot(snapshot configread.CurrentSnapshot) error { } binding, err := json.Marshal(struct { Task any `json:"task"` - Providers map[string]configread.CurrentProvider `json:"providers"` + Providers map[string]configread.Provider `json:"providers"` }{canonicalTask, providers}) if err != nil { return fmt.Errorf("encode immutable task binding: %w", err) } full, err := json.Marshal(struct { Task any `json:"task"` - Providers map[string]configread.CurrentProvider `json:"providers"` - SIP configread.CurrentSIP `json:"sip"` - Quota configread.CurrentQuota `json:"quota"` + Providers map[string]configread.Provider `json:"providers"` + SIP configread.SIP `json:"sip"` + Quota configread.Quota `json:"quota"` }{canonicalTask, providers, snapshot.SIP, snapshot.Quota}) if err != nil { return fmt.Errorf("encode task execution snapshot: %w", err) @@ -448,7 +448,7 @@ func (s *Store) SaveSnapshot(snapshot configread.CurrentSnapshot) error { // validateGlobalRevisions prevents two task bindings from silently // disagreeing about the same approved Dispatcher SIP or tenant quota revision. // Older revisions also cannot replace a newer binding across tasks. -func validateGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnapshot) error { +func validateGlobalRevisions(tx *sql.Tx, snapshot configread.Snapshot) error { newSIP, err := json.Marshal(snapshot.SIP) if err != nil { return fmt.Errorf("encode approved SIP revision: %w", err) @@ -470,8 +470,8 @@ func validateGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnapshot) er return fmt.Errorf("read global revision binding: %w", err) } var old struct { - SIP configread.CurrentSIP `json:"sip"` - Quota configread.CurrentQuota `json:"quota"` + SIP configread.SIP `json:"sip"` + Quota configread.Quota `json:"quota"` } if err := json.Unmarshal(raw, &old); err != nil { return fmt.Errorf("decode global revision binding for task %q: %w", taskID, err) @@ -515,30 +515,30 @@ func validateGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnapshot) er // ReadSnapshot returns only the current durable binding. Corrupt or mismatched // data is an error, never a signal to fetch an old contract instead. -func (s *Store) ReadSnapshot(dispatcherID string, tenantID int64, taskID string) (configread.CurrentSnapshot, error) { +func (s *Store) ReadSnapshot(dispatcherID string, tenantID int64, taskID string) (configread.Snapshot, error) { var body []byte var taskRevision, sipRevision, quotaRevision int64 err := s.db.QueryRow(`SELECT task_revision,sip_revision,quota_revision,snapshot_json FROM dispatcher_configs WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, dispatcherID, tenantID, taskID).Scan(&taskRevision, &sipRevision, "aRevision, &body) if err != nil { - return configread.CurrentSnapshot{}, fmt.Errorf("load task execution snapshot: %w", err) + return configread.Snapshot{}, fmt.Errorf("load task execution snapshot: %w", err) } var decoded struct { Task json.RawMessage `json:"task"` - Providers map[string]configread.CurrentProvider `json:"providers"` - SIP configread.CurrentSIP `json:"sip"` - Quota configread.CurrentQuota `json:"quota"` + Providers map[string]configread.Provider `json:"providers"` + SIP configread.SIP `json:"sip"` + Quota configread.Quota `json:"quota"` } if err := json.Unmarshal(body, &decoded); err != nil { - return configread.CurrentSnapshot{}, fmt.Errorf("decode task execution snapshot: %w", err) + return configread.Snapshot{}, fmt.Errorf("decode task execution snapshot: %w", err) } - var task configread.CurrentTask + var task configread.Task if err := json.Unmarshal(decoded.Task, &task); err != nil { - return configread.CurrentSnapshot{}, fmt.Errorf("decode stored task configuration: %w", err) + return configread.Snapshot{}, fmt.Errorf("decode stored task configuration: %w", err) } if task.DispatcherID != dispatcherID || task.TenantID != tenantID || task.TaskID != taskID || task.TaskRevision != taskRevision || decoded.SIP.DispatcherID != dispatcherID || decoded.SIP.Revision != sipRevision || decoded.Quota.DispatcherID != dispatcherID || decoded.Quota.TenantID != tenantID || decoded.Quota.QuotaRevision != quotaRevision { - return configread.CurrentSnapshot{}, errors.New("persisted task snapshot identity or revision mismatch") + return configread.Snapshot{}, errors.New("persisted task snapshot identity or revision mismatch") } - return configread.CurrentSnapshot{Task: task, Providers: decoded.Providers, SIP: decoded.SIP, Quota: decoded.Quota}, nil + return configread.Snapshot{Task: task, Providers: decoded.Providers, SIP: decoded.SIP, Quota: decoded.Quota}, nil } func (s *Store) ApplyControl(dispatcherID string, tenantID int64, taskID, action string) error { diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 8845e15..7e1f281 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -13,7 +13,7 @@ import ( const currentDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" -func currentStoreSnapshot(t *testing.T) configread.CurrentSnapshot { +func currentStoreSnapshot(t *testing.T) configread.Snapshot { t.Helper() dir := filepath.Join("..", "..", "contracts", "local", "examples") read := func(name string, dst any) { @@ -26,15 +26,15 @@ func currentStoreSnapshot(t *testing.T) configread.CurrentSnapshot { t.Fatal(err) } } - var s configread.CurrentSnapshot + var s configread.Snapshot read("config-read-task-asr", &s.Task) read("config-read-sip", &s.SIP) read("config-read-quota", &s.Quota) var p struct { - Providers []configread.CurrentProvider `json:"providers"` + Providers []configread.Provider `json:"providers"` } read("config-read-providers", &p) - s.Providers = make(map[string]configread.CurrentProvider) + s.Providers = make(map[string]configread.Provider) for _, provider := range p.Providers { s.Providers[provider.ProviderRef] = provider } @@ -47,8 +47,8 @@ func TestStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { if err != nil { t.Fatal(err) } - task := configread.CurrentDiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: 1, Status: "running"} - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{task}); err != nil { + task := configread.DiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: 1, Status: "running"} + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{task}); err != nil { t.Fatal(err) } if allowed, _ := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); allowed { @@ -77,7 +77,7 @@ func TestStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { if allowed, _ := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); allowed { t.Fatal("restart erased durable pause") } - if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.CurrentDiscoveredTask{task}); err != nil { + if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{task}); err != nil { t.Fatal(err) } if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil {