fix: InvalidClientTokenId panic + VT defaults + Dockerfile config + hardcore tests (58/61)

- Add InvalidClientTokenId, ValidationError to SqsErrors (was causing panic: WriteHeader code 0)
- Extract applyEnvironmentDefaults() with defer in LoadYamlConfig (VT=30 even without config file)
- Dockerfile: copy goaws.yaml to /conf/, pass --config flag
- Fix SendMessageBatch form parsing for AWS Query Protocol
- Add hardcore_test.sh (61 tests, 14 groups)
- Remove temp debug logs from delete_message.go
This commit is contained in:
Naeel
2026-04-09 16:11:32 +03:00
parent e5b1845acf
commit 6c1aadd886
6 changed files with 764 additions and 76 deletions
+26 -24
View File
@@ -23,6 +23,9 @@ var envs map[string]models.Environment
func LoadYamlConfig(filename string, env string) []string {
ports := []string{"4100"}
// Гарантируем что дефолты всегда применяются, даже если конфиг не найден
defer applyEnvironmentDefaults()
if filename == "" {
root, _ := filepath.Abs(".")
err := filepath.WalkDir(root, func(path string, d fs.DirEntry, err error) error {
@@ -77,30 +80,7 @@ models.LogFile = envs[env].LogFile
}
}
if models.CurrentEnvironment.QueueAttributeDefaults.VisibilityTimeout <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.VisibilityTimeout = 30
}
if models.CurrentEnvironment.QueueAttributeDefaults.MaximumMessageSize <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.MaximumMessageSize = 262144 // 256K
}
if models.CurrentEnvironment.QueueAttributeDefaults.MessageRetentionPeriod <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.MessageRetentionPeriod = 345600 // 4 days
}
if models.CurrentEnvironment.QueueAttributeDefaults.ReceiveMessageWaitTimeSeconds <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.ReceiveMessageWaitTimeSeconds = 0
}
if models.CurrentEnvironment.AccountID == "" {
models.CurrentEnvironment.AccountID = "queue"
}
if models.CurrentEnvironment.Host == "" {
models.CurrentEnvironment.Host = "localhost"
models.CurrentEnvironment.Port = "4100"
}
// Дефолты применяются через defer applyEnvironmentDefaults() в начале функции
models.SyncQueues.Lock()
for _, queue := range envs[env].Queues {
@@ -156,6 +136,28 @@ models.SyncQueues.Unlock()
return ports
}
// applyEnvironmentDefaults — применяет дефолтные значения для QueueAttributeDefaults,
// AccountID и Host. Вызывается через defer в LoadYamlConfig, чтобы дефолты
// устанавливались при любом раннем return (например, если конфиг не найден).
func applyEnvironmentDefaults() {
if models.CurrentEnvironment.QueueAttributeDefaults.VisibilityTimeout <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.VisibilityTimeout = 30
}
if models.CurrentEnvironment.QueueAttributeDefaults.MaximumMessageSize <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.MaximumMessageSize = 262144 // 256K
}
if models.CurrentEnvironment.QueueAttributeDefaults.MessageRetentionPeriod <= 0 {
models.CurrentEnvironment.QueueAttributeDefaults.MessageRetentionPeriod = 345600 // 4 days
}
if models.CurrentEnvironment.AccountID == "" {
models.CurrentEnvironment.AccountID = "queue"
}
if models.CurrentEnvironment.Host == "" {
models.CurrentEnvironment.Host = "localhost"
models.CurrentEnvironment.Port = "4100"
}
}
func setQueueRedrivePolicy(queues map[string]*models.Queue, q *models.Queue, strRedrivePolicy string) error {
// Поддерживаем maxReceiveCount как int и как string (AWS SDK использует string)
redrivePolicy1 := struct {
+50 -49
View File
@@ -3,62 +3,63 @@
package gosqs
import (
"net/http"
"strings"
"net/http"
"strings"
"shared-sqs/app/interfaces"
"shared-sqs/app/models"
"shared-sqs/app/utils"
"github.com/gorilla/mux"
log "github.com/sirupsen/logrus"
"shared-sqs/app/interfaces"
"shared-sqs/app/models"
"shared-sqs/app/utils"
"github.com/gorilla/mux"
log "github.com/sirupsen/logrus"
)
func DeleteMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewDeleteMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok {
log.Error("Invalid Request - DeleteMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
requestBody := models.NewDeleteMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok {
log.Error("Invalid Request - DeleteMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
t := getTenantFromContext(req)
if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
}
t := getTenantFromContext(req)
if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
}
receiptHandle := requestBody.ReceiptHandle
queueUrl := requestBody.QueueUrl
queueName := ""
if queueUrl == "" {
vars := mux.Vars(req)
queueName = vars["queueName"]
} else {
uriSegments := strings.Split(queueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
}
receiptHandle := requestBody.ReceiptHandle
queueUrl := requestBody.QueueUrl
queueName := ""
if queueUrl == "" {
vars := mux.Vars(req)
queueName = vars["queueName"]
} else {
uriSegments := strings.Split(queueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
}
key := tenantQueueKey(t.AccessKey, queueName)
log.Info("Deleting Message, Queue:", queueName, ", ReceiptHandle:", receiptHandle)
key := tenantQueueKey(t.AccessKey, queueName)
log.Info("Deleting Message, Queue:", queueName, ", ReceiptHandle:", receiptHandle)
models.SyncQueues.Lock()
defer models.SyncQueues.Unlock()
if _, ok := models.SyncQueues.Queues[key]; ok {
for i, msg := range models.SyncQueues.Queues[key].Messages {
if msg.ReceiptHandle == receiptHandle {
models.SyncQueues.Queues[key].UnlockGroup(msg.GroupID)
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages[:i], models.SyncQueues.Queues[key].Messages[i+1:]...)
delete(models.SyncQueues.Queues[key].Duplicates, msg.DeduplicationID)
respStruct := models.DeleteMessageResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return 200, &respStruct
}
}
log.Warning("Receipt Handle not found")
} else {
log.Warning("Queue not found")
}
models.SyncQueues.Lock()
defer models.SyncQueues.Unlock()
if _, ok := models.SyncQueues.Queues[key]; ok {
for i, msg := range models.SyncQueues.Queues[key].Messages {
if msg.ReceiptHandle == receiptHandle {
models.SyncQueues.Queues[key].UnlockGroup(msg.GroupID)
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages[:i], models.SyncQueues.Queues[key].Messages[i+1:]...)
delete(models.SyncQueues.Queues[key].Duplicates, msg.DeduplicationID)
respStruct := models.DeleteMessageResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return 200, &respStruct
}
}
log.Warning("Receipt Handle not found")
} else {
log.Warning("Queue not found")
}
return utils.CreateErrorResponseV1("MessageDoesNotExist", true)
return utils.CreateErrorResponseV1("MessageDoesNotExist", true)
}
+6
View File
@@ -16,6 +16,12 @@ func init() {
"MessageTooBig": {HttpError: http.StatusBadRequest, Type: "MessageTooBig", Code: "InvalidParameterValue", Message: "The message size exceeds the limit."},
"InvalidParameterValue": {HttpError: http.StatusBadRequest, Type: "InvalidParameterValue", Code: "AWS.SimpleQueueService.InvalidParameterValue", Message: "An invalid or out-of-range value was supplied for the input parameter."},
"InvalidAttributeValue": {HttpError: http.StatusBadRequest, Type: "InvalidAttributeValue", Code: "AWS.SimpleQueueService.InvalidAttributeValue", Message: "Invalid Value for the parameter RedrivePolicy."},
// InvalidClientTokenId — невалидные credentials тенанта
"InvalidClientTokenId": {HttpError: http.StatusForbidden, Type: "InvalidClientTokenId", Code: "AWS.SimpleQueueService.InvalidClientTokenId", Message: "The security token included in the request is invalid."},
// ValidationError — ошибка валидации параметров (например, VisibilityTimeout вне диапазона)
"ValidationError": {HttpError: http.StatusBadRequest, Type: "ValidationError", Code: "AWS.SimpleQueueService.ValidationError", Message: "The input fails to satisfy the constraints specified by an AWS service."},
// LimitExceeded — превышен лимит очередей тенанта (max_queues)
"LimitExceeded": {HttpError: http.StatusBadRequest, Type: "LimitExceeded", Code: "AWS.SimpleQueueService.LimitExceeded", Message: "You've reached the limit on the number of queues."},
}
SnsErrors = map[string]SnsErrorType{
"InvalidParameterValue": {HttpError: http.StatusBadRequest, Type: "InvalidParameterValue", Code: "AWS.SimpleNotificationService.InvalidParameterValue", Message: "An invalid or out-of-range value was supplied for the input parameter."},
+19 -2
View File
@@ -230,8 +230,25 @@ type SendMessageBatchRequest struct {
}
func (r *SendMessageBatchRequest) SetAttributesFromForm(values url.Values) {
for entryIndex := range r.Entries {
r.Entries[entryIndex].MessageAttributes = parseMessageAttributes(values, fmt.Sprintf("Entries.%d.MessageAttributes", entryIndex))
// Парсим записи по AWS Query Protocol: SendMessageBatchRequestEntry.N.Id (1-based)
// Gorilla/schema с дефолтными тегами ищет Entries.0.Id, что не соответствует AWS SQS API.
for i := 1; ; i++ {
id := values.Get(fmt.Sprintf("SendMessageBatchRequestEntry.%d.Id", i))
if id == "" {
break
}
entry := SendMessageBatchRequestEntry{
Id: id,
MessageBody: values.Get(fmt.Sprintf("SendMessageBatchRequestEntry.%d.MessageBody", i)),
MessageDeduplicationId: values.Get(fmt.Sprintf("SendMessageBatchRequestEntry.%d.MessageDeduplicationId", i)),
MessageGroupId: values.Get(fmt.Sprintf("SendMessageBatchRequestEntry.%d.MessageGroupId", i)),
}
ds := values.Get(fmt.Sprintf("SendMessageBatchRequestEntry.%d.DelaySeconds", i))
if ds != "" {
entry.DelaySeconds, _ = strconv.Atoi(ds)
}
entry.MessageAttributes = parseMessageAttributes(values, fmt.Sprintf("SendMessageBatchRequestEntry.%d.MessageAttribute", i))
r.Entries = append(r.Entries, entry)
}
}