Files
go-sip/internal/configread/snapshots.go
T

196 lines
6.7 KiB
Go

package configread
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/url"
"strconv"
"git.ipao.vip/rogee/go-sip/internal/contract"
)
// Snapshot carries the four independently validated configuration
// resources. Agent.Raw preserves explicit zero/false and immutable AI fields.
type Snapshot struct {
SIP SIP
Task Task
Quota Quota
Providers map[string]Provider
}
type SIP struct {
Resource string `json:"resource"`
DispatcherID string `json:"dispatcher_id"`
Revision int64 `json:"revision"`
Trunks json.RawMessage `json:"trunks"`
}
type Provider struct {
ProviderRef string `json:"provider_ref"`
Role string `json:"role"`
Enabled bool `json:"enabled"`
Adapter string `json:"adapter"`
Endpoint string `json:"endpoint"`
Credential string `json:"credential"`
}
type providerResponse struct {
Resource string `json:"resource"`
DispatcherID string `json:"dispatcher_id"`
Providers []Provider `json:"providers"`
}
type Agent struct {
Mode string `json:"mode"`
ASR struct {
ProviderRef string `json:"provider_ref"`
} `json:"asr"`
LLM struct {
ProviderRef string `json:"provider_ref"`
} `json:"llm"`
TTS struct {
ProviderRef string `json:"provider_ref"`
} `json:"tts"`
Raw json.RawMessage `json:"-"`
}
func (a *Agent) UnmarshalJSON(raw []byte) error {
type plain Agent
var parsed plain
if err := json.Unmarshal(raw, &parsed); err != nil {
return err
}
*a = Agent(parsed)
a.Raw = append([]byte(nil), raw...)
return nil
}
type Task struct {
Resource string `json:"resource"`
DispatcherID string `json:"dispatcher_id"`
TenantID int64 `json:"tenant_id"`
TaskID string `json:"task_id"`
TaskRevision int64 `json:"task_revision"`
Status string `json:"status"`
Name string `json:"name"`
MaxConcurrentCalls int64 `json:"max_concurrent_calls"`
RingTimeoutMS int64 `json:"ring_timeout_ms"`
MaxCallDurationMS int64 `json:"max_call_duration_ms"`
RoutePolicyID string `json:"route_policy_id"`
CallerProfileID string `json:"caller_profile_id"`
AllowedTrunkIDs []string `json:"allowed_trunk_ids"`
Schedule json.RawMessage `json:"schedule"`
Agent Agent `json:"agent"`
Raw json.RawMessage `json:"-"`
}
func (t *Task) UnmarshalJSON(raw []byte) error {
type plain Task
var parsed plain
if err := json.Unmarshal(raw, &parsed); err != nil {
return err
}
*t = Task(parsed)
t.Raw = append([]byte(nil), raw...)
return nil
}
type Quota struct {
Resource string `json:"resource"`
DispatcherID string `json:"dispatcher_id"`
TenantID int64 `json:"tenant_id"`
QuotaRevision int64 `json:"quota_revision"`
MaxConcurrentCalls int64 `json:"max_concurrent_calls"`
}
// 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) 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 SIP{}, errors.New("SIP configuration dispatcher owner mismatch")
}
return sip, nil
}
// ReadProviders reads the complete provider catalog once per Dispatcher boot.
// Running tasks use this immutable catalog until the next restart.
func (c *Client) ReadProviders(ctx context.Context) (map[string]Provider, error) {
var response providerResponse
if err := c.readResource(ctx, configReadPath+"/ai-providers", "ai_providers", &response); err != nil {
return nil, fmt.Errorf("read AI providers: %w", err)
}
if response.DispatcherID != c.dispatcherID {
return nil, errors.New("AI providers dispatcher owner mismatch")
}
providers := make(map[string]Provider, len(response.Providers))
for _, provider := range response.Providers {
if _, exists := providers[provider.ProviderRef]; exists {
return nil, fmt.Errorf("duplicate AI provider reference %q", provider.ProviderRef)
}
providers[provider.ProviderRef] = provider
}
return providers, nil
}
// ReadTask reads current task and quota resources using the SIP and provider
// snapshots approved at boot. No task operation fetches the global resources.
func (c *Client) ReadTask(ctx context.Context, taskID string, tenantID int64, sip SIP, providers map[string]Provider) (Snapshot, error) {
if taskID == "" || tenantID <= 0 {
return Snapshot{}, errors.New("task ID and positive tenant ID are required")
}
if sip.DispatcherID != c.dispatcherID || sip.Revision <= 0 || len(providers) == 0 {
return Snapshot{}, errors.New("task requires approved SIP and AI provider snapshots")
}
result := Snapshot{SIP: sip, Providers: providers}
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.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 || result.Task.DispatcherID != c.dispatcherID || result.Quota.DispatcherID != c.dispatcherID || result.Task.TenantID != tenantID || result.Quota.TenantID != tenantID || result.Task.TaskID != taskID {
return Snapshot{}, errors.New("configuration dispatcher, task, or tenant owner mismatch")
}
for _, expected := range []struct{ ref, role string }{
{result.Task.Agent.ASR.ProviderRef, "asr"},
{result.Task.Agent.LLM.ProviderRef, "llm"},
{result.Task.Agent.TTS.ProviderRef, "tts"},
} {
if expected.ref == "" {
continue
}
provider, found := result.Providers[expected.ref]
if !found || !provider.Enabled || provider.Role != expected.role {
return Snapshot{}, fmt.Errorf("AI provider %q is missing, disabled, or has an incompatible role", expected.ref)
}
}
return result, nil
}
func (c *Client) readResource(ctx context.Context, path, resource string, dst any) error {
body, err := c.get(ctx, path)
if err != nil {
return err
}
if err := contract.ValidateCurrent("config-read", body); err != nil {
return err
}
var identity struct {
Resource string `json:"resource"`
}
if err := json.Unmarshal(body, &identity); err != nil {
return err
}
if identity.Resource != resource {
return fmt.Errorf("expected %s, got %s", resource, identity.Resource)
}
return json.Unmarshal(body, dst)
}