feat(shared-sqs): queue CRUD + message peek/send/purge in UI (v0.1.6)
- admin.go: 5 new endpoints: create/delete queue, peek/send/purge messages - index.html: expandable message rows, send modal, msg detail modal, create queue modal - deployment.yaml: update image to naeel/shared-sqs:v0.1.6 - Docker image pushed: naeel/shared-sqs:v0.1.6
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
// app/admin/admin.go
|
||||
// Admin API handlers for shared-sqs management
|
||||
// Created: 2026-04-09
|
||||
// Updated: 2026-04-10 — добавлены endpoints для управления очередями и просмотра сообщений
|
||||
package admin
|
||||
|
||||
import (
|
||||
@@ -12,10 +13,24 @@ import (
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/tenant"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/gorilla/mux"
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
// ─── вспомогательная функция: найти очередь тенанта по имени ───────────────
|
||||
// findQueue — возвращает ключ и очередь тенанта по имени, или "",nil если не найдено.
|
||||
func findQueue(tenantAccessKey, queueName string) (string, *models.Queue) {
|
||||
key := tenantAccessKey + ":" + queueName
|
||||
models.SyncQueues.RLock()
|
||||
q, ok := models.SyncQueues.Queues[key]
|
||||
models.SyncQueues.RUnlock()
|
||||
if !ok {
|
||||
return "", nil
|
||||
}
|
||||
return key, q
|
||||
}
|
||||
|
||||
// Handler — admin API handler, holds TenantStore и admin token
|
||||
type Handler struct {
|
||||
store *tenant.TenantStore
|
||||
@@ -27,7 +42,7 @@ func NewHandler(store *tenant.TenantStore, adminToken string) *Handler {
|
||||
return &Handler{store: store, adminToken: adminToken}
|
||||
}
|
||||
|
||||
// RegisterRoutes — регистрирует admin маршруты на переданном router
|
||||
// RegisterRoutes — регистрирует admin маршруты на переданном router (с bearer auth)
|
||||
func (h *Handler) RegisterRoutes(r *mux.Router) {
|
||||
adminRouter := r.PathPrefix("/admin").Subrouter()
|
||||
adminRouter.Use(h.bearerAuthMiddleware)
|
||||
@@ -36,11 +51,17 @@ func (h *Handler) RegisterRoutes(r *mux.Router) {
|
||||
adminRouter.HandleFunc("/tenants/{id}", h.getTenant).Methods("GET")
|
||||
adminRouter.HandleFunc("/tenants/{id}", h.deleteTenant).Methods("DELETE")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues", h.listTenantQueues).Methods("GET")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues", h.createTenantQueue).Methods("POST")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues/{queue}", h.deleteTenantQueue).Methods("DELETE")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.peekQueueMessages).Methods("GET")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.sendMessageToQueue).Methods("POST")
|
||||
adminRouter.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.purgeQueue).Methods("DELETE")
|
||||
adminRouter.HandleFunc("/health", h.detailedHealth).Methods("GET")
|
||||
}
|
||||
|
||||
// RegisterPublicRoutes — публичные маршруты для UI console (без auth)
|
||||
// Дублируют admin API, но доступны без bearer token для удобства демо
|
||||
// TODO: убрать или заменить на session-auth перед production
|
||||
func (h *Handler) RegisterPublicRoutes(r *mux.Router) {
|
||||
ui := r.PathPrefix("/ui/api").Subrouter()
|
||||
ui.HandleFunc("/health", h.detailedHealth).Methods("GET")
|
||||
@@ -49,6 +70,11 @@ func (h *Handler) RegisterPublicRoutes(r *mux.Router) {
|
||||
ui.HandleFunc("/tenants/{id}", h.getTenant).Methods("GET")
|
||||
ui.HandleFunc("/tenants/{id}", h.deleteTenant).Methods("DELETE")
|
||||
ui.HandleFunc("/tenants/{id}/queues", h.listTenantQueues).Methods("GET")
|
||||
ui.HandleFunc("/tenants/{id}/queues", h.createTenantQueue).Methods("POST")
|
||||
ui.HandleFunc("/tenants/{id}/queues/{queue}", h.deleteTenantQueue).Methods("DELETE")
|
||||
ui.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.peekQueueMessages).Methods("GET")
|
||||
ui.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.sendMessageToQueue).Methods("POST")
|
||||
ui.HandleFunc("/tenants/{id}/queues/{queue}/messages", h.purgeQueue).Methods("DELETE")
|
||||
}
|
||||
|
||||
// bearerAuthMiddleware — проверяет Bearer token для admin API (Trap #12)
|
||||
@@ -245,6 +271,216 @@ func (h *Handler) listTenantQueues(w http.ResponseWriter, r *http.Request) {
|
||||
json.NewEncoder(w).Encode(queues)
|
||||
}
|
||||
|
||||
// ─── QUEUE MANAGEMENT HANDLERS ────────────────────────────────────────────
|
||||
|
||||
// createQueueRequest — тело запроса POST .../queues
|
||||
type createQueueRequest struct {
|
||||
Name string `json:"name"`
|
||||
}
|
||||
|
||||
// createTenantQueue — POST /admin/tenants/{id}/queues
|
||||
// Создаёт новую очередь для тенанта прямо в SyncQueues (без SQS-протокола).
|
||||
// Проверяет лимит MaxQueues тенанта и уникальность имени.
|
||||
func (h *Handler) createTenantQueue(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
tid := vars["id"]
|
||||
t, ok := h.store.GetByID(tid)
|
||||
if !ok {
|
||||
jsonErr(w, http.StatusNotFound, "tenant not found")
|
||||
return
|
||||
}
|
||||
var req createQueueRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Name == "" {
|
||||
jsonErr(w, http.StatusBadRequest, "name is required")
|
||||
return
|
||||
}
|
||||
key := t.AccessKey + ":" + req.Name
|
||||
models.SyncQueues.Lock()
|
||||
if _, exists := models.SyncQueues.Queues[key]; exists {
|
||||
models.SyncQueues.Unlock()
|
||||
jsonErr(w, http.StatusConflict, "queue already exists")
|
||||
return
|
||||
}
|
||||
// Проверяем лимит очередей тенанта
|
||||
count := 0
|
||||
for k := range models.SyncQueues.Queues {
|
||||
if strings.HasPrefix(k, t.AccessKey+":") {
|
||||
count++
|
||||
}
|
||||
}
|
||||
if t.MaxQueues > 0 && count >= t.MaxQueues {
|
||||
models.SyncQueues.Unlock()
|
||||
jsonErr(w, http.StatusForbidden, "queue limit exceeded")
|
||||
return
|
||||
}
|
||||
models.SyncQueues.Queues[key] = &models.Queue{
|
||||
Name: req.Name,
|
||||
VisibilityTimeout: 30,
|
||||
MaximumMessageSize: 262144,
|
||||
MessageRetentionPeriod: 345600,
|
||||
Messages: []models.SqsMessage{},
|
||||
Duplicates: make(map[string]time.Time),
|
||||
}
|
||||
models.SyncQueues.Unlock()
|
||||
log.Infof("admin: created queue %s for tenant %s", req.Name, t.ID)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusCreated)
|
||||
json.NewEncoder(w).Encode(map[string]string{"name": req.Name, "status": "created"})
|
||||
}
|
||||
|
||||
// deleteTenantQueue — DELETE /admin/tenants/{id}/queues/{queue}
|
||||
// Удаляет очередь тенанта из SyncQueues вместе со всеми её сообщениями.
|
||||
func (h *Handler) deleteTenantQueue(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
tid, queueName := vars["id"], vars["queue"]
|
||||
t, ok := h.store.GetByID(tid)
|
||||
if !ok {
|
||||
jsonErr(w, http.StatusNotFound, "tenant not found")
|
||||
return
|
||||
}
|
||||
key := t.AccessKey + ":" + queueName
|
||||
models.SyncQueues.Lock()
|
||||
if _, exists := models.SyncQueues.Queues[key]; !exists {
|
||||
models.SyncQueues.Unlock()
|
||||
jsonErr(w, http.StatusNotFound, "queue not found")
|
||||
return
|
||||
}
|
||||
delete(models.SyncQueues.Queues, key)
|
||||
models.SyncQueues.Unlock()
|
||||
log.Infof("admin: deleted queue %s for tenant %s", queueName, t.ID)
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// peekMessageItem — одно сообщение в ответе peekQueueMessages (без receipt handle)
|
||||
type peekMessageItem struct {
|
||||
ID string `json:"id"`
|
||||
Body string `json:"body"`
|
||||
MD5 string `json:"md5"`
|
||||
SentAt string `json:"sent_at"`
|
||||
Receives int `json:"receives"`
|
||||
InFlight bool `json:"in_flight"`
|
||||
}
|
||||
|
||||
// peekQueueMessages — GET /admin/tenants/{id}/queues/{queue}/messages?limit=50
|
||||
// Peek-просмотр сообщений: НЕ удаляет, НЕ выставляет ReceiptHandle — только чтение.
|
||||
// Это принципиальное отличие от SQS ReceiveMessage (который скрывает сообщения).
|
||||
func (h *Handler) peekQueueMessages(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
tid, queueName := vars["id"], vars["queue"]
|
||||
t, ok := h.store.GetByID(tid)
|
||||
if !ok {
|
||||
jsonErr(w, http.StatusNotFound, "tenant not found")
|
||||
return
|
||||
}
|
||||
_, q := findQueue(t.AccessKey, queueName)
|
||||
if q == nil {
|
||||
jsonErr(w, http.StatusNotFound, "queue not found")
|
||||
return
|
||||
}
|
||||
// Лимит по умолчанию 50, максимум 1000
|
||||
limit := 50
|
||||
if lv := r.URL.Query().Get("limit"); lv != "" {
|
||||
if n := 0; len(lv) > 0 {
|
||||
for _, c := range lv {
|
||||
if c < '0' || c > '9' {
|
||||
n = -1
|
||||
break
|
||||
}
|
||||
n = n*10 + int(c-'0')
|
||||
}
|
||||
if n > 0 && n <= 1000 {
|
||||
limit = n
|
||||
}
|
||||
}
|
||||
}
|
||||
models.SyncQueues.RLock()
|
||||
result := make([]peekMessageItem, 0, len(q.Messages))
|
||||
for i, msg := range q.Messages {
|
||||
if i >= limit {
|
||||
break
|
||||
}
|
||||
result = append(result, peekMessageItem{
|
||||
ID: msg.Uuid,
|
||||
Body: msg.MessageBody,
|
||||
MD5: msg.MD5OfMessageBody,
|
||||
SentAt: msg.SentTime.Format(time.RFC3339),
|
||||
Receives: msg.NumberOfReceives,
|
||||
InFlight: msg.ReceiptHandle != "",
|
||||
})
|
||||
}
|
||||
models.SyncQueues.RUnlock()
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(result)
|
||||
}
|
||||
|
||||
// sendMessageRequest — тело запроса POST .../messages
|
||||
type sendMessageRequest struct {
|
||||
Body string `json:"body"`
|
||||
}
|
||||
|
||||
// sendMessageToQueue — POST /admin/tenants/{id}/queues/{queue}/messages
|
||||
// Отправляет сообщение напрямую в очередь минуя SQS-протокол.
|
||||
// Используется только из UI console — для prod нужен нормальный SQS send.
|
||||
func (h *Handler) sendMessageToQueue(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
tid, queueName := vars["id"], vars["queue"]
|
||||
t, ok := h.store.GetByID(tid)
|
||||
if !ok {
|
||||
jsonErr(w, http.StatusNotFound, "tenant not found")
|
||||
return
|
||||
}
|
||||
key, q := findQueue(t.AccessKey, queueName)
|
||||
if q == nil {
|
||||
jsonErr(w, http.StatusNotFound, "queue not found")
|
||||
return
|
||||
}
|
||||
var req sendMessageRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Body == "" {
|
||||
jsonErr(w, http.StatusBadRequest, "body is required")
|
||||
return
|
||||
}
|
||||
msg := models.SqsMessage{
|
||||
MessageBody: req.Body,
|
||||
Uuid: uuid.NewString(),
|
||||
SentTime: time.Now(),
|
||||
}
|
||||
models.SyncQueues.Lock()
|
||||
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
|
||||
models.SyncQueues.Unlock()
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusCreated)
|
||||
json.NewEncoder(w).Encode(map[string]string{"id": msg.Uuid, "status": "sent"})
|
||||
}
|
||||
|
||||
// purgeQueue — DELETE /admin/tenants/{id}/queues/{queue}/messages
|
||||
// Удаляет все сообщения из очереди (purge). Сама очередь остаётся.
|
||||
func (h *Handler) purgeQueue(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
tid, queueName := vars["id"], vars["queue"]
|
||||
t, ok := h.store.GetByID(tid)
|
||||
if !ok {
|
||||
jsonErr(w, http.StatusNotFound, "tenant not found")
|
||||
return
|
||||
}
|
||||
key, q := findQueue(t.AccessKey, queueName)
|
||||
if q == nil {
|
||||
jsonErr(w, http.StatusNotFound, "queue not found")
|
||||
return
|
||||
}
|
||||
models.SyncQueues.Lock()
|
||||
models.SyncQueues.Queues[key].Messages = models.SyncQueues.Queues[key].Messages[:0]
|
||||
models.SyncQueues.Unlock()
|
||||
log.Infof("admin: purged queue %s for tenant %s", queueName, t.ID)
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// jsonErr — вспомогательная функция: ответ с ошибкой в JSON
|
||||
func jsonErr(w http.ResponseWriter, code int, msg string) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(code)
|
||||
json.NewEncoder(w).Encode(map[string]string{"error": msg})
|
||||
}
|
||||
|
||||
// adminHealthDetail — ответ GET /admin/health
|
||||
type adminHealthDetail struct {
|
||||
Status string `json:"status"`
|
||||
|
||||
Reference in New Issue
Block a user