From 6a41bbd51961975829bce3f6d84a0543acdd8ea9 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 14:52:53 +0800 Subject: [PATCH] refactor(contracts): keep historical local schemas out of runtime --- contracts/contracts.go | 12 +- contracts/current_read_test.go | 26 +++ contracts/local_v04_test.go | 56 ------ .../saas-dispatcher-implementation.md | 1 + internal/contract/command_test.go | 32 --- internal/contract/contract.go | 187 ------------------ internal/contract/contract_test.go | 132 ------------- internal/contract/current.go | 3 + internal/contract/event_mq.go | 38 ---- internal/contract/event_mq_test.go | 47 ----- .../contract/local_mock_recording_failure.go | 32 --- .../local_mock_recording_failure_test.go | 34 ---- internal/contract/schema.go | 131 ------------ internal/contract/schema_source_test.go | 23 --- internal/contract/schema_test.go | 123 ------------ internal/contract/service.go | 86 -------- internal/contract/service_test.go | 39 ---- 17 files changed, 32 insertions(+), 970 deletions(-) delete mode 100644 contracts/local_v04_test.go delete mode 100644 internal/contract/command_test.go delete mode 100644 internal/contract/contract.go delete mode 100644 internal/contract/contract_test.go delete mode 100644 internal/contract/event_mq.go delete mode 100644 internal/contract/event_mq_test.go delete mode 100644 internal/contract/local_mock_recording_failure.go delete mode 100644 internal/contract/local_mock_recording_failure_test.go delete mode 100644 internal/contract/schema.go delete mode 100644 internal/contract/schema_source_test.go delete mode 100644 internal/contract/schema_test.go delete mode 100644 internal/contract/service.go delete mode 100644 internal/contract/service_test.go diff --git a/contracts/contracts.go b/contracts/contracts.go index 563fae2..09b6d68 100644 --- a/contracts/contracts.go +++ b/contracts/contracts.go @@ -5,13 +5,12 @@ import ( "encoding/json" "fmt" "path" - "strings" ) -// Files contains pinned upstream schemas and the active local contract versions. +// Files contains the pinned upstream history and only the active local contract. // Runtime code never reads a checkout or resolves schema refs online. // -//go:embed upstream local/v0.1 local/v0.2 local/v0.3 local/v0.4 local/*.json +//go:embed upstream local/*.json local/examples var Files embed.FS const SourceCommit = "v1" @@ -27,13 +26,6 @@ func ReadCurrent(name string) ([]byte, error) { } } -func ReadLocal(version, name string) ([]byte, error) { - if (version != "v0.1" && version != "v0.2" && version != "v0.3" && version != "v0.4") || name == "" || path.Base(name) != name || !strings.HasSuffix(name, ".schema.json") { - return nil, fmt.Errorf("invalid project-local schema %q/%q", version, name) - } - return Files.ReadFile(path.Join("local", version, name)) -} - func Read(name string) ([]byte, error) { return Files.ReadFile(path.Join("upstream", SourceCommit, name)) } diff --git a/contracts/current_read_test.go b/contracts/current_read_test.go index 3ce98a1..92148b4 100644 --- a/contracts/current_read_test.go +++ b/contracts/current_read_test.go @@ -7,6 +7,32 @@ import ( "git.ipao.vip/rogee/go-sip/contracts" ) +func TestPinnedUpstreamBundleRemainsReadable(t *testing.T) { + if body, err := contracts.Read("mq.schema.json"); err != nil || len(body) == 0 { + t.Fatalf("pinned upstream schema missing: bytes=%d err=%v", len(body), err) + } + var schema any + if err := contracts.ReadJSON("mq.schema.json", &schema); err != nil || schema == nil { + t.Fatalf("pinned upstream schema is invalid: %v", err) + } + if _, err := contracts.Read("missing-contract.json"); err == nil { + t.Fatal("missing upstream source was accepted") + } +} + +func TestRuntimeBundleExcludesHistoricalLocalContracts(t *testing.T) { + for _, name := range []string{ + "local/v0.1/call-result-v0.1-proposal.schema.json", + "local/v0.2/config-read-v0.2.schema.json", + "local/v0.3/config-read-v0.3.schema.json", + "local/v0.4/call-execute-v0.4-proposal.schema.json", + } { + if _, err := contracts.Files.ReadFile(name); err == nil { + t.Fatalf("historical contract %s is still embedded in the active runtime", name) + } + } +} + func TestCurrentContractRead(t *testing.T) { for _, name := range []string{"config-read.schema.json", "task-discovery.schema.json", "mq.schema.json", "mq-topology.json", "manifest.json"} { data, err := contracts.ReadCurrent(name) diff --git a/contracts/local_v04_test.go b/contracts/local_v04_test.go deleted file mode 100644 index 9169a18..0000000 --- a/contracts/local_v04_test.go +++ /dev/null @@ -1,56 +0,0 @@ -package contracts - -import ( - "bytes" - "encoding/json" - "testing" -) - -func TestReadLocalV04SchemasAndRejectInvalidPaths(t *testing.T) { - for _, name := range []string{ - "task-discovery-v0.4-proposal.schema.json", - "task-control-v0.4-proposal.schema.json", - "call-execute-v0.4-proposal.schema.json", - } { - data, err := ReadLocal("v0.4", name) - if err != nil || !json.Valid(data) { - t.Fatalf("embedded v0.4 schema %s: valid=%v err=%v", name, json.Valid(data), err) - } - if name == "task-discovery-v0.4-proposal.schema.json" && bytes.Contains(data, []byte("task-discovery-v0.3-proposal.schema.json")) { - t.Fatal("v0.4 task discovery must not reference the historical v0.3 schema") - } - original, err := Files.ReadFile("local/v0.4/" + name) - if err != nil || !bytes.Equal(data, original) { - t.Fatalf("runtime schema differs from embedded bundle: %s err=%v", name, err) - } - } - for _, input := range []struct{ version, name string }{ - {"v0.2", "task-control-v0.4-proposal.schema.json"}, - {"v0.4", ""}, - {"v0.4", "../task-control-v0.4-proposal.schema.json"}, - {"v0.4", "task-control-v0.4-proposal.json"}, - {"v0.4", "does-not-exist.schema.json"}, - } { - if _, err := ReadLocal(input.version, input.name); err == nil { - t.Fatalf("invalid or missing embedded schema accepted: %+v", input) - } - } -} - -func TestReadUpstreamContractAndDecodeErrors(t *testing.T) { - data, err := Read("mq.schema.json") - if err != nil || !json.Valid(data) { - t.Fatalf("upstream MQ schema: valid=%v err=%v", json.Valid(data), err) - } - var schema map[string]json.RawMessage - if err := ReadJSON("mq.schema.json", &schema); err != nil || len(schema) == 0 { - t.Fatalf("decode upstream MQ schema: fields=%d err=%v", len(schema), err) - } - if err := ReadJSON("missing.schema.json", &schema); err == nil { - t.Fatal("missing upstream contract did not return an error") - } - var scalar int - if err := ReadJSON("mq.schema.json", &scalar); err == nil { - t.Fatal("invalid upstream schema destination did not return a decode error") - } -} diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 19f5b74..cef1f4d 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -126,6 +126,7 @@ - 旧 ARI/媒体执行包:`internal/callruntime` 仅实现旧 `ai.Snapshot` 的独立 `Run` 入口,无当前 Agent/Dispatcher 调用者;删除该包及仅针对该路径的测试。现行获批执行、录音恢复和完整 AI Mock 测试继续覆盖唯一现行入口;删除死代码不表示真实 Asterisk/ARI 或线路已验证,历史验收报告保留当时包名与覆盖率事实。 - 旧通话流程入口:旧 `Execute`/`ExecuteWithCapture` 使用可缺省的旧 AI 快照,现已无 Agent 调用者;删除入口及专属测试,保留同一 `executeFlow` 媒体顺序实现。先在当前 `ExecuteApproved` 的测试中补齐首次媒体等待上限、三轮对话、开场发送失败不伪造播放、拒联后不发送回答与 ASR-only 实际采集区间,再迁移共用测试夹具;现行获批通话测试和全仓测试通过。未声称真实媒体链路通过。 - 旧 AI 配置与供应商链路:旧 `Snapshot`/独立 MQ 配置响应/独立授权、旧可变 Mock 与旧供应商流水线已无现行执行者;删除专属实现和测试。当前获批执行仍使用原值 AI 模式及每通话不可变参数;将仍被当前路径使用的 `Mode` 与 `TurnResult` 移至 `internal/ai/turn.go`。先在当前 SDK 隔离测试补齐 LLM 失败与空回答候选不自动重试,再删除旧 SDK 专属测试;ASR-only、完整 AI、关键词、时限和媒体的当前测试继续通过。旧供应商测试的历史通过不能代替真实供应商联调。 +- 旧项目内合同装载:测试先复现 `contracts.Files` 将 v0.1–v0.4 的历史本地 Schema 仍嵌入现行二进制;现在仅嵌入获批当前 local Schema/示例和原样上游 v1 来源,移除可从代码访问旧 local 的 `ReadLocal`、旧命令/事件/配置验证层及专属测试,将并发 Schema 缓存保留于当前验证器。项目历史 Schema 文件与外部上游来源仍原样留在仓库,没有改写来源哈希;历史文件不再构成运行时回退。当前合同正反例、上游来源读取和全仓 Go 测试通过,未代表外部 SaaS 已签收。 ## 验收台账 diff --git a/internal/contract/command_test.go b/internal/contract/command_test.go deleted file mode 100644 index 51792c7..0000000 --- a/internal/contract/command_test.go +++ /dev/null @@ -1,32 +0,0 @@ -package contract - -import ( - "encoding/json" - "git.ipao.vip/rogee/go-sip/contracts" - "testing" -) - -func TestDecodeMQCommandV1(t *testing.T) { - for _, name := range []string{"call-execute", "task-control", "call-replay", "command-replay"} { - raw, err := contracts.Files.ReadFile("upstream/" + MQSourceCommit + "/examples/" + name + ".json") - if err != nil { - t.Fatal(err) - } - message, err := DecodeMQCommand(raw) - if err != nil { - t.Fatal(err) - } - if message.DispatcherID == "" || message.CommandType == "" { - t.Fatal("command identity lost") - } - var changed map[string]any - if err := json.Unmarshal(raw, &changed); err != nil { - t.Fatal(err) - } - changed["schema_version"] = "1.0" - bad, _ := json.Marshal(changed) - if _, err := DecodeMQCommand(bad); err == nil { - t.Fatal("legacy command accepted") - } - } -} diff --git a/internal/contract/contract.go b/internal/contract/contract.go deleted file mode 100644 index a8ef019..0000000 --- a/internal/contract/contract.go +++ /dev/null @@ -1,187 +0,0 @@ -package contract - -import ( - "bytes" - "encoding/json" - "errors" - "fmt" - "strings" - "time" - "unicode/utf8" - - "git.ipao.vip/rogee/go-sip/contracts" - "github.com/santhosh-tekuri/jsonschema/v6" -) - -type CommandEnvelope struct { - DispatcherID string `json:"dispatcher_id,omitempty"` - SchemaVersion string `json:"schema_version"` - CommandType string `json:"command_type"` - CommandID string `json:"command_id"` - TenantID string `json:"tenant_id"` - TenantKey string `json:"tenant_key"` - TraceID string `json:"trace_id"` - IssuedAt string `json:"issued_at"` - NotAfter string `json:"not_after"` - Payload json.RawMessage `json:"payload"` -} - -type ExecutePayload struct { - ExecutionID string `json:"execution_id"` - TaskID string `json:"task_id"` - TaskItemID string `json:"task_item_id"` - TaskRevision int64 `json:"task_revision"` - Callee string `json:"callee"` - RoutePolicyID string `json:"route_policy_id"` - CallerProfileID string `json:"caller_profile_id"` - AgentVersionID string `json:"agent_version_id"` - Variables map[string]any `json:"variables"` - RingTimeoutMS int64 `json:"ring_timeout_ms"` - MaxCallDurationMS int64 `json:"max_call_duration_ms"` -} - -type EventEnvelope struct { - SchemaVersion string `json:"schema_version"` - EventID string `json:"event_id"` - EventType string `json:"event_type"` - TenantID string `json:"tenant_id"` - TenantKey string `json:"tenant_key"` - TraceID string `json:"trace_id"` - OccurredAt string `json:"occurred_at"` - AggregateType string `json:"aggregate_type"` - AggregateID string `json:"aggregate_id"` - AggregateVersion int64 `json:"aggregate_version"` - Payload map[string]any `json:"payload"` -} - -var ErrInvalidTenantKey = errors.New("invalid tenant_key") - -// ValidateJSON applies the imported JSON Schema. It intentionally validates -// the source contract instead of maintaining a second hand-written schema. -func ValidateJSON(raw []byte) error { - return ValidateSourceSchema("mq.schema.json", raw) -} - -// ValidateEvent applies the event-specific payload contract in addition to the -// generic MQ envelope contract. -func ValidateEvent(raw []byte) error { - return ValidateSourceSchema("event-payloads.schema.json", raw) -} - -func ValidateLocalConfigRead(raw []byte) error { - return validateLocalSchema("config-read-v0.3.schema.json", raw) -} - -func ValidateLocalTaskDiscoveryV04(raw []byte) error { - return validateLocalSchema("task-discovery-v0.4-proposal.schema.json", raw) -} - -func ValidateLocalAIAuthorization(raw []byte) error { - return validateLocalSchema("ai-authorization-v0.2.schema.json", raw) -} - -func ValidateLocalTaskControlV04(raw []byte) error { - return validateLocalSchema("task-control-v0.4-proposal.schema.json", raw) -} - -func ValidateLocalCallExecuteV04(raw []byte) error { - return validateLocalSchema("call-execute-v0.4-proposal.schema.json", raw) -} - -func ValidateLocalCommandNext(raw []byte) error { - return validateLocalSchema("command-next-v0.1-proposal.schema.json", raw) -} - -func ValidateLocalCallResult(raw []byte) error { - return validateLocalSchema("call-result-v0.1-proposal.schema.json", raw) -} - -func validateLocalSchema(schemaName string, raw []byte) error { - value, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw)) - if err != nil { - return fmt.Errorf("decode %s response: %w", schemaName, err) - } - schema, err := localSchemaFromBundle(schemaName) - if err != nil { - return err - } - if err := schema.Validate(value); err != nil { - var validation *jsonschema.ValidationError - if errors.As(err, &validation) { - return fmt.Errorf("%s validation failed at /%s", schemaName, strings.Join(validation.InstanceLocation, "/")) - } - return fmt.Errorf("%s validation failed", schemaName) - } - return nil -} - -func ValidateSourceSchema(schemaName string, raw []byte) error { - value, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw)) - if err != nil { - return fmt.Errorf("decode json: %w", err) - } - schema, err := schemaFromBundle(contracts.SourceCommit, schemaName) - if err != nil { - return err - } - if err := schema.Validate(value); err != nil { - return fmt.Errorf("%s validation: %w", schemaName, err) - } - return nil -} - -func DecodeExecute(raw []byte) (CommandEnvelope, ExecutePayload, error) { - envelope, err := DecodeMQCommand(raw) - if err != nil { - return CommandEnvelope{}, ExecutePayload{}, err - } - if envelope.CommandType != "call.execute" { - return CommandEnvelope{}, ExecutePayload{}, fmt.Errorf("unsupported command_type %q", envelope.CommandType) - } - if err := ValidateTenantKey(envelope.TenantKey); err != nil { - return CommandEnvelope{}, ExecutePayload{}, err - } - var payload ExecutePayload - if err := json.Unmarshal(envelope.Payload, &payload); err != nil { - return CommandEnvelope{}, ExecutePayload{}, fmt.Errorf("decode call.execute payload: %w", err) - } - return envelope, payload, nil -} - -func ValidateTenantKey(key string) error { - if key == "" || !utf8.ValidString(key) { - return fmt.Errorf("%w: must be non-empty valid UTF-8", ErrInvalidTenantKey) - } - if len([]byte(key)) > 224 { - return fmt.Errorf("%w: %d UTF-8 bytes exceeds 224-byte routing budget", ErrInvalidTenantKey, len([]byte(key))) - } - return nil -} - -func NotAfterExpired(raw string, now time.Time) (bool, error) { - deadline, err := time.Parse(time.RFC3339Nano, raw) - if err != nil { - return false, fmt.Errorf("parse not_after: %w", err) - } - return !now.Before(deadline), nil -} - -func CloneJSON(raw []byte) json.RawMessage { - return bytes.Clone(raw) -} - -type EventBuilder struct { - DispatcherID string - TenantID string - TenantKey string - TraceID string - EventType string - Aggregate string - AggregateID string - Version int64 - Payload map[string]any -} - -func (b EventBuilder) Marshal(now time.Time, eventID string) ([]byte, error) { - return b.MarshalMQ(b.DispatcherID, now, eventID) -} diff --git a/internal/contract/contract_test.go b/internal/contract/contract_test.go deleted file mode 100644 index 2385827..0000000 --- a/internal/contract/contract_test.go +++ /dev/null @@ -1,132 +0,0 @@ -package contract - -import ( - "strings" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/contracts" - "git.ipao.vip/rogee/go-sip/internal/testfixture" -) - -func TestDecodeExecuteValidatesImportedSchema(t *testing.T) { - raw, err := testfixture.Execute() - if err != nil { - t.Fatal(err) - } - envelope, payload, err := DecodeExecute(raw) - if err != nil { - t.Fatal(err) - } - if envelope.TenantKey != "tenant-demo-key" || payload.ExecutionID == "" { - t.Fatalf("unexpected decoded command: %+v %+v", envelope, payload) - } -} - -func TestValidateTenantKeyPreservesRawValueAndBudget(t *testing.T) { - for _, key := range []string{"客户-A", " tenant /raw "} { - if err := ValidateTenantKey(key); err != nil { - t.Fatalf("tenant key %q: %v", key, err) - } - } - if err := ValidateTenantKey(strings.Repeat("x", 225)); err == nil { - t.Fatal("expected routing budget rejection") - } -} - -func TestEventFixtures(t *testing.T) { - positive := []string{ - "examples/event-command-result.json", - "examples/event-call-status.json", - "examples/event-transcript-updated.json", - "examples/event-call-finished.json", - "examples/event-recording-uploaded.json", - "examples/event-recording-failed.json", - "examples/event-transcript-failed.json", - "examples/event-contact-opt-out.json", - } - for _, name := range positive { - raw, err := contracts.Read(name) - if err != nil { - t.Fatalf("%s: %v", name, err) - } - if err := ValidateEvent(raw); err != nil { - t.Fatalf("%s: %v", name, err) - } - } - invalid, err := contracts.Read("examples/invalid-v1.json") - if err != nil { - t.Fatal(err) - } - if err := ValidateEvent(invalid); err == nil { - t.Fatal("expected invalid transcript alias to be rejected") - } -} - -func TestProjectOwnedW01Fixtures(t *testing.T) { - valid := []string{ - "examples/ai-authorization.json", - "examples/oss-upload-grant.json", - "examples/static-cell-artifact.json", - "p1-development-profile.json", - } - for _, name := range valid { - raw, err := contracts.Read(name) - if err != nil { - t.Fatalf("%s: %v", name, err) - } - if err := ValidateSourceSchema(schemaNameForFixture(name), raw); err != nil { - t.Fatalf("%s: %v", name, err) - } - } - invalid := []struct { - name string - schema string - }{ - {"examples/invalid-ai-authorization-revoked.json", "ai-authorization.schema.json"}, - {"examples/invalid-oss-upload-http.json", "oss-upload.schema.json"}, - } - for _, tc := range invalid { - raw, err := contracts.Read(tc.name) - if err != nil { - t.Fatalf("%s: %v", tc.name, err) - } - if err := ValidateSourceSchema(tc.schema, raw); err == nil { - t.Fatalf("%s: expected rejection", tc.name) - } - } -} - -func schemaNameForFixture(name string) string { - switch name { - case "examples/ai-authorization.json": - return "ai-authorization.schema.json" - case "examples/oss-upload-grant.json": - return "oss-upload.schema.json" - case "examples/static-cell-artifact.json": - return "static-cell-artifact.schema.json" - case "p1-development-profile.json": - return "p1-development-profile.schema.json" - default: - return "" - } -} - -func TestEventBuilder(t *testing.T) { - raw, err := (EventBuilder{ - DispatcherID: "11111111-1111-4111-8111-111111111111", - TenantID: "tenant-1", TenantKey: "tenant-1", TraceID: "trace-1", - EventType: "command.result", Aggregate: "command", AggregateID: "cmd-1", Version: 1, - Payload: map[string]any{ - "command_id": "cmd-1", "command_type": "call.execute", - "status": "accepted", "reason_code": "accepted", - "execution_id": "execution-1", - }, - }).Marshal(time.Unix(0, 0), "event-1") - if err != nil { - t.Fatal(err) - } - if err := ValidateEvent(raw); err != nil { - t.Fatal(err) - } -} diff --git a/internal/contract/current.go b/internal/contract/current.go index 2fb82af..b8de87b 100644 --- a/internal/contract/current.go +++ b/internal/contract/current.go @@ -5,11 +5,14 @@ import ( "errors" "fmt" "strings" + "sync" "git.ipao.vip/rogee/go-sip/contracts" "github.com/santhosh-tekuri/jsonschema/v6" ) +var compiledSchemas sync.Map + // ValidateCurrent checks only the current project-local HTTP/MQ contract. // An absent schema is an error; it never falls back to historical bundles. func ValidateCurrent(surface string, raw []byte) error { diff --git a/internal/contract/event_mq.go b/internal/contract/event_mq.go deleted file mode 100644 index 7e4ca98..0000000 --- a/internal/contract/event_mq.go +++ /dev/null @@ -1,38 +0,0 @@ -package contract - -import ( - "encoding/json" - "fmt" - "time" -) - -// MarshalMQ constructs the Dispatcher-owned external event envelope. Internal -// Agent event frames do not supply the Dispatcher identity or select a version. -func (b EventBuilder) MarshalMQ(dispatcherID string, now time.Time, eventID string) ([]byte, error) { - if err := ValidateTenantKey(b.TenantKey); err != nil { - return nil, err - } - if b.Version < 1 { - return nil, fmt.Errorf("aggregate version must be positive") - } - envelope := struct { - EventEnvelope - DispatcherID string `json:"dispatcher_id"` - }{EventEnvelope: EventEnvelope{ - SchemaVersion: "2.0", EventID: eventID, EventType: b.EventType, - TenantID: b.TenantID, TenantKey: b.TenantKey, TraceID: b.TraceID, - OccurredAt: now.UTC().Format(time.RFC3339Nano), AggregateType: b.Aggregate, - AggregateID: b.AggregateID, AggregateVersion: b.Version, Payload: b.Payload, - }, DispatcherID: dispatcherID} - raw, err := json.Marshal(envelope) - if err != nil { - return nil, err - } - if err := ValidateMQMessage(raw); err != nil { - return nil, err - } - if err := ValidateEvent(raw); err != nil { - return nil, err - } - return raw, nil -} diff --git a/internal/contract/event_mq_test.go b/internal/contract/event_mq_test.go deleted file mode 100644 index 06a6165..0000000 --- a/internal/contract/event_mq_test.go +++ /dev/null @@ -1,47 +0,0 @@ -package contract - -import ( - "encoding/json" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/contracts" -) - -func TestMQEventBuilderUsesStrictDispatcherEnvelope(t *testing.T) { - raw, err := contracts.Files.ReadFile("upstream/" + MQSourceCommit + "/examples/event-command-result.json") - if err != nil { - t.Fatal(err) - } - var event EventEnvelope - if err := json.Unmarshal(raw, &event); err != nil { - t.Fatal(err) - } - builder := EventBuilder{TenantID: event.TenantID, TenantKey: event.TenantKey, TraceID: event.TraceID, EventType: event.EventType, Aggregate: event.AggregateType, AggregateID: event.AggregateID, Version: event.AggregateVersion, Payload: event.Payload} - now, err := time.Parse(time.RFC3339Nano, event.OccurredAt) - if err != nil { - t.Fatal(err) - } - const id = "11111111-1111-4111-8111-111111111111" - result, err := builder.MarshalMQ(id, now, event.EventID) - if err != nil { - t.Fatal(err) - } - if err := ValidateMQMessage(result); err != nil { - t.Fatal(err) - } - var values map[string]any - if err := json.Unmarshal(result, &values); err != nil { - t.Fatal(err) - } - if values["schema_version"] != "2.0" || values["dispatcher_id"] != id { - t.Fatal("missing v2 identity") - } - if _, err := builder.MarshalMQ("", now, event.EventID); err == nil { - t.Fatal("missing dispatcher accepted") - } - builder.Payload["unexpected"] = true - if _, err := builder.MarshalMQ(id, now, event.EventID); err == nil { - t.Fatal("unknown payload field accepted") - } -} diff --git a/internal/contract/local_mock_recording_failure.go b/internal/contract/local_mock_recording_failure.go deleted file mode 100644 index 9c9c86f..0000000 --- a/internal/contract/local_mock_recording_failure.go +++ /dev/null @@ -1,32 +0,0 @@ -package contract - -import ( - "encoding/json" - "fmt" -) - -const LocalMockRecordingFailureVersion = "local-mock-recording-failure.v0.1" -const MaxLocalMockRecordingFailureBytes = 4096 - -// LocalMockRecordingFailure is the strict, Mock-only payload of an Agent -// recording.progress failure fact. It is not a SaaS event or a real-mode API. -type LocalMockRecordingFailure struct { - SchemaVersion string `json:"schema_version"` - UploadID string `json:"upload_id"` - RecordingID string `json:"recording_id"` - ErrorCode string `json:"error_code"` -} - -func DecodeLocalMockRecordingFailure(raw []byte) (LocalMockRecordingFailure, error) { - if len(raw) == 0 || len(raw) > MaxLocalMockRecordingFailureBytes { - return LocalMockRecordingFailure{}, fmt.Errorf("local Mock recording failure payload must be 1–%d bytes", MaxLocalMockRecordingFailureBytes) - } - if err := validateLocalSchema("local-mock-recording-failure-v0.1.schema.json", raw); err != nil { - return LocalMockRecordingFailure{}, err - } - var fact LocalMockRecordingFailure - if err := json.Unmarshal(raw, &fact); err != nil { - return LocalMockRecordingFailure{}, fmt.Errorf("decode local Mock recording failure: %w", err) - } - return fact, nil -} diff --git a/internal/contract/local_mock_recording_failure_test.go b/internal/contract/local_mock_recording_failure_test.go deleted file mode 100644 index 6fb46c5..0000000 --- a/internal/contract/local_mock_recording_failure_test.go +++ /dev/null @@ -1,34 +0,0 @@ -package contract - -import ( - "strings" - "testing" -) - -func TestLocalMockRecordingFailureSchema(t *testing.T) { - for _, code := range []string{"upload_authorization_failed", "upload_authorization_expired", "upload_failed", "checksum_mismatch"} { - t.Run(code, func(t *testing.T) { - raw := []byte(`{"schema_version":"local-mock-recording-failure.v0.1","upload_id":"upload-a","recording_id":"recording-a","error_code":"` + code + `"}`) - fact, err := DecodeLocalMockRecordingFailure(raw) - if err != nil || fact.UploadID != "upload-a" || fact.RecordingID != "recording-a" || fact.ErrorCode != code { - t.Fatalf("valid local Mock failure: %+v err=%v", fact, err) - } - }) - } - for _, tc := range []struct { - name string - raw string - }{ - {"missing upload", `{"schema_version":"local-mock-recording-failure.v0.1","recording_id":"recording-a","error_code":"upload_failed"}`}, - {"unknown field", `{"schema_version":"local-mock-recording-failure.v0.1","upload_id":"upload-a","recording_id":"recording-a","error_code":"upload_failed","put_url":"https://example.invalid/secret"}`}, - {"timeout owned by Dispatcher", `{"schema_version":"local-mock-recording-failure.v0.1","upload_id":"upload-a","recording_id":"recording-a","error_code":"upload_timeout"}`}, - {"wrong version", `{"schema_version":"v0.2","upload_id":"upload-a","recording_id":"recording-a","error_code":"upload_failed"}`}, - {"oversized", `{"schema_version":"local-mock-recording-failure.v0.1","upload_id":"upload-a","recording_id":"recording-a","error_code":"upload_failed"}` + strings.Repeat(" ", 4097)}, - } { - t.Run(tc.name, func(t *testing.T) { - if _, err := DecodeLocalMockRecordingFailure([]byte(tc.raw)); err == nil { - t.Fatal("unversioned, oversized or incompatible Mock failure accepted") - } - }) - } -} diff --git a/internal/contract/schema.go b/internal/contract/schema.go deleted file mode 100644 index a53104b..0000000 --- a/internal/contract/schema.go +++ /dev/null @@ -1,131 +0,0 @@ -package contract - -import ( - "bytes" - "fmt" - "path" - "strings" - "sync" - - "git.ipao.vip/rogee/go-sip/contracts" - "github.com/santhosh-tekuri/jsonschema/v6" -) - -var compiledSchemas sync.Map - -type bundleLoader struct{ version string } - -func (l bundleLoader) base() string { - return "https://go-sip.local/contracts/" + l.version + "/" -} - -// Load never accesses the network or resolves references against another bundle. -func (l bundleLoader) Load(address string) (any, error) { - if !strings.HasPrefix(address, l.base()) { - return nil, fmt.Errorf("schema reference outside pinned bundle: %s", address) - } - name := strings.TrimPrefix(address, l.base()) - if name == "" || path.Base(name) != name || !strings.HasSuffix(name, ".schema.json") { - return nil, fmt.Errorf("invalid schema resource name %q", name) - } - raw, err := contracts.Files.ReadFile("upstream/" + l.version + "/" + name) - if err != nil { - return nil, fmt.Errorf("read pinned schema %s: %w", name, err) - } - return jsonschema.UnmarshalJSON(bytes.NewReader(raw)) -} - -const projectLocalSchemaBase = "https://go-sip.local/contracts/proposals/" - -type projectBundleLoader struct{} - -func (projectBundleLoader) Load(address string) (any, error) { - upstream := bundleLoader{version: contracts.SourceCommit} - if strings.HasPrefix(address, upstream.base()) { - return upstream.Load(address) - } - if !strings.HasPrefix(address, projectLocalSchemaBase) { - return nil, fmt.Errorf("schema reference outside pinned project bundles: %s", address) - } - name := strings.TrimPrefix(address, projectLocalSchemaBase) - if name == "" || path.Base(name) != name || !strings.HasSuffix(name, ".schema.json") { - return nil, fmt.Errorf("invalid project-local schema resource %q", name) - } - raw, err := contracts.ReadLocal(localSchemaVersion(name), name) - if err != nil { - return nil, fmt.Errorf("read project-local schema %s: %w", name, err) - } - return jsonschema.UnmarshalJSON(bytes.NewReader(raw)) -} - -func localSchemaVersion(name string) string { - switch name { - case "config-read-v0.1.schema.json", "command-next-v0.1-proposal.schema.json", "call-result-v0.1-proposal.schema.json", "local-mock-recording-failure-v0.1.schema.json": - return "v0.1" - case "config-read-v0.2.schema.json": - return "v0.2" - case "config-read-v0.3.schema.json", "ai-authorization-v0.2.schema.json", "static-cell-artifact-v0.2.schema.json", "task-discovery-v0.3-proposal.schema.json": - return "v0.3" - case "task-discovery-v0.4-proposal.schema.json", "task-control-v0.4-proposal.schema.json", "call-execute-v0.4-proposal.schema.json": - return "v0.4" - default: - return "" - } -} - -func localSchemaFromBundle(name string) (*jsonschema.Schema, error) { - version := localSchemaVersion(name) - if version == "" { - return nil, fmt.Errorf("unknown project-local schema %q", name) - } - key := "local/" + version + "/" + name - if cached, ok := compiledSchemas.Load(key); ok { - return cached.(*jsonschema.Schema), nil - } - raw, err := contracts.ReadLocal(version, name) - if err != nil { - return nil, err - } - doc, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw)) - if err != nil { - return nil, fmt.Errorf("decode project-local schema %s: %w", name, err) - } - resource := projectLocalSchemaBase + name - compiler := jsonschema.NewCompiler() - compiler.UseLoader(projectBundleLoader{}) - compiler.AssertFormat() - if err := compiler.AddResource(resource, doc); err != nil { - return nil, fmt.Errorf("register project-local schema %s: %w", name, err) - } - schema, err := compiler.Compile(resource) - if err != nil { - return nil, fmt.Errorf("compile project-local schema %s: %w", name, err) - } - actual, _ := compiledSchemas.LoadOrStore(key, schema) - return actual.(*jsonschema.Schema), nil -} - -func schemaFromBundle(version, name string) (*jsonschema.Schema, error) { - key := version + "/" + name - if cached, ok := compiledSchemas.Load(key); ok { - return cached.(*jsonschema.Schema), nil - } - loader := bundleLoader{version: version} - resource := loader.base() + name - doc, err := loader.Load(resource) - if err != nil { - return nil, err - } - compiler := jsonschema.NewCompiler() - compiler.UseLoader(loader) - compiler.AssertFormat() - if err := compiler.AddResource(resource, doc); err != nil { - return nil, fmt.Errorf("register %s: %w", name, err) - } - schema, err := compiler.Compile(resource) - if err != nil { - return nil, fmt.Errorf("compile %s: %w", name, err) - } - actual, _ := compiledSchemas.LoadOrStore(key, schema) - return actual.(*jsonschema.Schema), nil -} diff --git a/internal/contract/schema_source_test.go b/internal/contract/schema_source_test.go deleted file mode 100644 index dc20332..0000000 --- a/internal/contract/schema_source_test.go +++ /dev/null @@ -1,23 +0,0 @@ -package contract - -import ( - "testing" - - "git.ipao.vip/rogee/go-sip/contracts" -) - -func TestContractBundleIsReadable(t *testing.T) { - for _, name := range []string{ - "mq.schema.json", "event-payloads.schema.json", "ai-config.schema.json", - "ai-authorization.schema.json", "oss-upload.schema.json", - "static-cell-artifact.schema.json", "p1-development-profile.schema.json", - } { - data, err := contracts.Read(name) - if err != nil { - t.Fatalf("%s: %v", name, err) - } - if len(data) == 0 { - t.Fatalf("%s is empty", name) - } - } -} diff --git a/internal/contract/schema_test.go b/internal/contract/schema_test.go deleted file mode 100644 index 08aabfa..0000000 --- a/internal/contract/schema_test.go +++ /dev/null @@ -1,123 +0,0 @@ -package contract - -import ( - "bytes" - "encoding/json" - "os" - "path/filepath" - "testing" - - "git.ipao.vip/rogee/go-sip/contracts" - "github.com/santhosh-tekuri/jsonschema/v6" -) - -func TestV1RuntimeSchemaResolvesOnlyEmbeddedContracts(t *testing.T) { - const version = contracts.SourceCommit - schema, err := schemaFromBundle(version, "mq.schema.json") - if err != nil { - t.Fatal(err) - } - for _, name := range []string{"command-query", "call-query-result", "ai-config-result", "call-execute", "event-recording-uploaded"} { - raw, err := contracts.Files.ReadFile("upstream/" + version + "/examples/" + name + ".json") - if err != nil { - t.Fatal(err) - } - value, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw)) - if err != nil { - t.Fatal(err) - } - if err := schema.Validate(value); err != nil { - t.Fatalf("%s: %v", name, err) - } - } - loader := bundleLoader{version: version} - for _, address := range []string{"https://example.invalid/mq.schema.json", "file:///etc/passwd", "https://go-sip.local/contracts/other/mq.schema.json"} { - if _, err := loader.Load(address); err == nil { - t.Fatalf("external/cross-version schema accepted: %s", address) - } - } - again, err := schemaFromBundle(version, "mq.schema.json") - if err != nil || again != schema { - t.Fatal("immutable compiled schema not reused") - } -} - -func TestBundleSchemaRejectsInvalidResources(t *testing.T) { - for _, name := range []string{"../mq.schema.json", "", "missing.schema.json"} { - if _, err := schemaFromBundle(contracts.SourceCommit, name); err == nil { - t.Fatalf("invalid resource accepted: %q", name) - } - } -} - -func TestProjectLocalConfigurationSchemasValidatePositivesAndRejectNegatives(t *testing.T) { - read := func(name string) []byte { - t.Helper() - body, err := os.ReadFile(filepath.Join("..", "..", "docs", "contracts", "examples", name)) - if err != nil { - t.Fatal(err) - } - return body - } - if err := ValidateLocalConfigRead(read("config-read-task-v0.1.json")); err != nil { - t.Fatalf("valid task config: %v", err) - } - var task map[string]any - if err := json.Unmarshal(read("config-read-task-v0.1.json"), &task); err != nil { - t.Fatal(err) - } - delete(task, "name") - delete(task, "group_id") - withoutDisplayFields, err := json.Marshal(task) - if err != nil { - t.Fatal(err) - } - if err := ValidateLocalConfigRead(withoutDisplayFields); err != nil { - t.Fatalf("task without optional display fields: %v", err) - } - task["name"] = 123 - withInvalidDisplayField, err := json.Marshal(task) - if err != nil { - t.Fatal(err) - } - if err := ValidateLocalConfigRead(withInvalidDisplayField); err == nil { - t.Fatal("invalid optional task name accepted") - } - if err := ValidateLocalConfigRead(read("config-read-invalid-extra-property-v0.1.json")); err == nil { - t.Fatal("config-read schema accepted an additional property") - } -} - -func TestV04LocalSchemasValidateExamples(t *testing.T) { - for _, tt := range []struct { - name string - validate func([]byte) error - valid bool - }{ - {"task-discovery-snapshot-page1-v0.4.json", ValidateLocalTaskDiscoveryV04, true}, - {"task-discovery-snapshot-page2-v0.4.json", ValidateLocalTaskDiscoveryV04, true}, - {"task-discovery-changes-v0.4.json", ValidateLocalTaskDiscoveryV04, true}, - {"task-discovery-error-v0.4.json", ValidateLocalTaskDiscoveryV04, true}, - {"task-discovery-invalid-v0.4.json", ValidateLocalTaskDiscoveryV04, false}, - {"task-control-pause-v0.4.json", ValidateLocalTaskControlV04, true}, - {"task-control-resume-v0.4.json", ValidateLocalTaskControlV04, true}, - {"task-control-stop-v0.4.json", ValidateLocalTaskControlV04, true}, - {"task-control-invalid-drain-v0.4.json", ValidateLocalTaskControlV04, false}, - {"task-control-invalid-expiry-v0.4.json", ValidateLocalTaskControlV04, false}, - {"call-execute-old-v0.4.json", ValidateLocalCallExecuteV04, true}, - {"call-execute-recent-v0.4.json", ValidateLocalCallExecuteV04, true}, - {"call-execute-altcallee-v0.4.json", ValidateLocalCallExecuteV04, true}, - {"call-execute-invalid-expiry-v0.4.json", ValidateLocalCallExecuteV04, false}, - } { - t.Run(tt.name, func(t *testing.T) { - body, err := os.ReadFile(filepath.Join("..", "..", "docs", "contracts", "examples", tt.name)) - if err != nil { - t.Fatal(err) - } - err = tt.validate(body) - if (err == nil) != tt.valid { - t.Fatalf("valid=%v, err=%v", tt.valid, err) - } - }) - } -} diff --git a/internal/contract/service.go b/internal/contract/service.go deleted file mode 100644 index 9a0d188..0000000 --- a/internal/contract/service.go +++ /dev/null @@ -1,86 +0,0 @@ -package contract - -import ( - "bytes" - "encoding/json" - "errors" - "fmt" - "unicode/utf8" - - "git.ipao.vip/rogee/go-sip/contracts" - "github.com/santhosh-tekuri/jsonschema/v6" -) - -// MQSourceCommit pins the approved MQ protocol. No incoming version selects a -// different decoder or causes fallback to the previous protocol. -const MQSourceCommit = contracts.SourceCommit - -var ErrInvalidServiceMessage = errors.New("invalid MQ service message") - -type ServiceMessage struct { - SchemaVersion string `json:"schema_version"` - MessageType string `json:"message_type"` - MessageID string `json:"message_id"` - DispatcherID string `json:"dispatcher_id"` - TenantID string `json:"tenant_id"` - TenantKey string `json:"tenant_key"` - TraceID string `json:"trace_id"` - IssuedAt string `json:"issued_at"` - NotAfter string `json:"not_after,omitempty"` - CorrelationID string `json:"correlation_id,omitempty"` - Status string `json:"status,omitempty"` - ReasonCode string `json:"reason_code,omitempty"` - Payload json.RawMessage `json:"payload"` -} - -func ValidateMQMessage(raw []byte) error { - if len(raw) == 0 || len(raw) > 256<<10 || !utf8.Valid(raw) { - return fmt.Errorf("%w: require UTF-8 JSON within 256 KiB", ErrInvalidServiceMessage) - } - value, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw)) - if err != nil { - return fmt.Errorf("%w: decode JSON: %w", ErrInvalidServiceMessage, err) - } - schema, err := schemaFromBundle(MQSourceCommit, "mq.schema.json") - if err != nil { - return err - } - if err := schema.Validate(value); err != nil { - return fmt.Errorf("%w: schema validation: %w", ErrInvalidServiceMessage, err) - } - return nil -} - -func DecodeMQCommand(raw []byte) (CommandEnvelope, error) { - if err := ValidateMQMessage(raw); err != nil { - return CommandEnvelope{}, err - } - var command CommandEnvelope - if err := json.Unmarshal(raw, &command); err != nil { - return CommandEnvelope{}, err - } - if command.CommandType == "" { - return CommandEnvelope{}, fmt.Errorf("%w: expected command", ErrInvalidServiceMessage) - } - if len([]byte(command.TenantKey)) > 196 { - return CommandEnvelope{}, ErrInvalidTenantKey - } - return command, nil -} - -func DecodeService(raw []byte) (ServiceMessage, error) { - if err := ValidateMQMessage(raw); err != nil { - return ServiceMessage{}, err - } - var message ServiceMessage - if err := json.Unmarshal(raw, &message); err != nil { - return ServiceMessage{}, err - } - if message.MessageType == "" { - return ServiceMessage{}, fmt.Errorf("%w: not a service request or response", ErrInvalidServiceMessage) - } - if len([]byte(message.TenantKey)) > 196 { - return ServiceMessage{}, fmt.Errorf("%w: exceeds 196 UTF-8 bytes", ErrInvalidTenantKey) - } - return message, nil -} diff --git a/internal/contract/service_test.go b/internal/contract/service_test.go deleted file mode 100644 index 23e8ad8..0000000 --- a/internal/contract/service_test.go +++ /dev/null @@ -1,39 +0,0 @@ -package contract - -import ( - "encoding/json" - "testing" - - "git.ipao.vip/rogee/go-sip/contracts" -) - -func TestDecodeServiceUsesApprovedMessageTypeAndCorrelation(t *testing.T) { - for _, name := range []string{"command-query", "call-query", "ai-config-request", "command-query-result", "call-query-result", "ai-config-result"} { - raw, err := contracts.Files.ReadFile("upstream/v1/examples/" + name + ".json") - if err != nil { - t.Fatal(err) - } - message, err := DecodeService(raw) - if err != nil { - t.Fatalf("%s: %v", name, err) - } - if message.MessageType == "" || message.DispatcherID == "" || message.MessageID == "" { - t.Fatal("message identity was lost") - } - var body map[string]any - if err := json.Unmarshal(raw, &body); err != nil { - t.Fatal(err) - } - body["message_kind"] = "request" - bad, _ := json.Marshal(body) - if _, err := DecodeService(bad); err == nil { - t.Fatal("invented message_kind accepted") - } - delete(body, "message_kind") - body["schema_version"] = "1.0" - bad, _ = json.Marshal(body) - if _, err := DecodeService(bad); err == nil { - t.Fatal("legacy service envelope accepted") - } - } -}