refactor(contracts): keep historical local schemas out of runtime

This commit is contained in:
2026-09-30 14:52:53 +08:00
parent 75b415deff
commit 6a41bbd519
17 changed files with 32 additions and 970 deletions
+2 -10
View File
@@ -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))
}
+26
View File
@@ -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)
-56
View File
@@ -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")
}
}
@@ -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 已签收。
## 验收台账
-32
View File
@@ -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")
}
}
}
-187
View File
@@ -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)
}
-132
View File
@@ -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)
}
}
+3
View File
@@ -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 {
-38
View File
@@ -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
}
-47
View File
@@ -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")
}
}
@@ -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
}
@@ -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")
}
})
}
}
-131
View File
@@ -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
}
-23
View File
@@ -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)
}
}
}
-123
View File
@@ -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)
}
})
}
}
-86
View File
@@ -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
}
-39
View File
@@ -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")
}
}
}