87 lines
2.8 KiB
Go
87 lines
2.8 KiB
Go
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
|
|
}
|