Files
go-sip/internal/dispatcher/consumer.go
T

92 lines
3.0 KiB
Go

package dispatcher
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"strings"
"git.ipao.vip/rogee/go-sip/internal/contract"
"git.ipao.vip/rogee/go-sip/internal/mq"
"git.ipao.vip/rogee/go-sip/internal/store"
)
// AcceptMQMessage returns only after the input decision and response are durable.
func (d *Dispatcher) AcceptMQMessage(body []byte, routingKey string) error {
var discriminator struct {
MessageType string `json:"message_type"`
CommandType string `json:"command_type"`
}
if err := json.Unmarshal(body, &discriminator); err != nil {
return mq.Permanent(fmt.Errorf("decode MQ message: %w", err))
}
var err error
switch discriminator.MessageType {
case "ai.config.result":
err = d.store.StoreAIConfigResponse(body, routingKey)
if err == nil {
slog.Info("MQ AI configuration response persisted")
}
case "command.query", "call.query":
var responseID string
var duplicate bool
responseID, duplicate, err = d.store.HandleQuery(body, routingKey)
if err == nil {
slog.Info("MQ query response persisted", "message_type", discriminator.MessageType, "response_id", responseID, "duplicate", duplicate)
}
case "":
switch discriminator.CommandType {
case "task.control":
var responseID string
var duplicate bool
d.executionMu.Lock()
responseID, duplicate, err = d.store.HandleTaskControl(body, routingKey)
d.executionMu.Unlock()
if err == nil {
slog.Info("MQ task control persisted", "response_id", responseID, "duplicate", duplicate)
}
case "call.replay", "command.replay":
var responseID string
var duplicate bool
responseID, duplicate, err = d.store.HandleReplay(body, routingKey)
if err == nil {
slog.Info("MQ replay decision persisted", "command_type", discriminator.CommandType, "response_id", responseID, "duplicate", duplicate)
}
default:
_, err = d.AcceptCommand(body, routingKey)
}
default:
return mq.Permanent(errors.New("unsupported inbound MQ service message"))
}
if err != nil && isPermanentCommandError(err) {
return mq.Permanent(err)
}
return err
}
func (d *Dispatcher) ConsumeTenant(ctx context.Context, broker *mq.Broker, tenantKey string) error {
if broker == nil {
return errors.New("broker is required")
}
if err := contract.ValidateTenantKey(tenantKey); err != nil {
return err
}
queue, err := broker.DeclareTenantQueue(tenantKey)
if err != nil {
return err
}
return broker.Consume(ctx, queue, func(ctx context.Context, routingKey string, body []byte) error {
return d.AcceptMQMessage(body, routingKey)
})
}
func isPermanentCommandError(err error) bool {
if errors.Is(err, contract.ErrInvalidServiceMessage) || errors.Is(err, contract.ErrInvalidTenantKey) || errors.Is(err, store.ErrMessageScope) || errors.Is(err, store.ErrIdempotencyConflict) || errors.Is(err, store.ErrAIConfigResponse) || errors.Is(err, store.ErrCommandConflict) {
return true
}
message := err.Error()
return strings.Contains(message, "schema validation") || strings.Contains(message, "decode json") || strings.Contains(message, "routing mismatch")
}