Sprint 7: race conditions + read_at + slog

1. chat.EnsureChat: ON CONFLICT DO UPDATE (race-safe)
2. chat.Send: tx.Begin/Commit (atomic INSERT message + UPDATE last_msg_at)
3. chat.MarkRead: добавил member-check (NOT a member -> 403)
4. chat.Message: +ReadAt field
5. migrations/0006: read_at column on messages
6. locations.Nearby: bounding-box prefilter + CTE (haversine только для отфильтрованных)
7. audit.Log: real metadata JSON passthrough (no more 'null' TODO)
8. main.go: slog JSON logger (был stdlib log)
9. WS: 'message_read' event от MarkRead
10. WS push: добавлены read/read_at/edited в message payload
11. Bumped v0.6.0
This commit is contained in:
ga 2026-08-20 21:04:56 +00:00
parent 2d3bdfb98f
commit 0f9e0ce596
7 changed files with 158 additions and 77 deletions

View File

@ -4,7 +4,7 @@ import (
"context"
"fmt"
"io"
"log"
"log/slog"
"os"
"os/signal"
"syscall"
@ -12,8 +12,8 @@ import (
"github.com/gofiber/fiber/v2"
"github.com/gofiber/fiber/v2/middleware/cors"
flog "github.com/gofiber/fiber/v2/middleware/logger"
"github.com/gofiber/fiber/v2/middleware/limiter"
"github.com/gofiber/fiber/v2/middleware/logger"
"github.com/gofiber/fiber/v2/middleware/recover"
"github.com/minio/minio-go/v7"
@ -35,40 +35,47 @@ import (
)
func main() {
// структурный логгер — JSON для prod
logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
slog.SetDefault(logger)
cfg, err := config.Load()
if err != nil {
log.Fatalf("config: %v", err)
logger.Error("config load failed", "err", err)
os.Exit(1)
}
ctx := context.Background()
pg, err := db.New(ctx, cfg.PostgresDSN())
if err != nil {
log.Fatalf("postgres: %v", err)
logger.Error("postgres connect failed", "err", err)
os.Exit(1)
}
defer pg.Close()
log.Println("postgres: connected")
logger.Info("postgres connected")
migDir := os.Getenv("MIGRATIONS_DIR")
if migDir == "" {
migDir = "./migrations"
}
if err := db.RunMigrations(ctx, pg, migDir); err != nil {
log.Fatalf("migrations: %v", err)
logger.Error("migrations failed", "err", err)
os.Exit(1)
}
rdb, err := redis.New(ctx, cfg.RedisAddr(), cfg.RedisPassword, cfg.RedisDB)
if err != nil {
log.Printf("redis warning: %v", err)
logger.Warn("redis not available", "err", err)
} else {
log.Println("redis: connected")
logger.Info("redis connected")
}
st, err := storage.New(ctx, cfg.MinIOEndpoint, cfg.MinIOAccessKey, cfg.MinIOSecretKey, cfg.MinIOBucket, cfg.MinIOUseSSL)
if err != nil {
log.Printf("minio warning: %v", err)
logger.Warn("minio not available", "err", err)
} else {
log.Println("minio: connected")
logger.Info("minio connected")
}
usersRepo := users.NewRepo(pg.Pool)
@ -87,9 +94,9 @@ func main() {
if cfg.TelegramBotToken != "" {
bot = telegram.NewBot(cfg.TelegramBotToken)
tgH = handlers.NewTGHandler(bot, usersRepo, authSvc, consentRepo, auditRepo, cfg.TelegramWebhookSecret)
log.Println("telegram bot: configured")
logger.Info("telegram bot configured")
} else {
log.Println("telegram bot: NOT configured (set TELEGRAM_BOT_TOKEN)")
logger.Warn("telegram bot NOT configured")
}
authH := handlers.NewAuthHandler(cfg, usersRepo, consentRepo, auditRepo, authSvc)
@ -107,7 +114,7 @@ func main() {
WriteTimeout: 15 * time.Second,
})
app.Use(recover.New())
app.Use(logger.New())
app.Use(flog.New())
app.Use(cors.New(cors.Config{
AllowOrigins: "https://buhapp.mygoodservice.ru,https://t.me,https://web.telegram.org",
AllowHeaders: "Origin, Content-Type, Accept, Authorization, X-Telegram-Bot-Api-Secret-Token",
@ -135,7 +142,7 @@ func main() {
return c.JSON(fiber.Map{
"status": "ok",
"time": time.Now().UTC().Format(time.RFC3339),
"version": "0.5.0",
"version": "0.6.0",
"db": pg.Pool != nil,
"redis": rdb != nil,
"storage": st != nil,
@ -152,7 +159,6 @@ func main() {
return c.JSON(fiber.Map{"version": cfg.DisclaimerVersion, "url": "/legal/DISCLAIMER.md"})
})
// публичная раздача MinIO через бэкенд (MinIO в docker network, не на хосте)
if st != nil {
app.Get("/storage/*", func(c *fiber.Ctx) error {
key := c.Params("*")
@ -224,7 +230,7 @@ func main() {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
<-sigCh
log.Println("shutdown...")
logger.Info("shutdown initiated")
_ = app.ShutdownWithTimeout(10 * time.Second)
}()
@ -232,13 +238,13 @@ func main() {
webhookURL := cfg.TelegramWebhookURL
if webhookURL != "" {
if err := bot.SetWebhook(webhookURL, cfg.TelegramWebhookSecret); err != nil {
log.Printf("telegram setWebhook failed: %v (fallback to long poll)", err)
logger.Warn("telegram setWebhook failed, will use long poll", "err", err)
} else {
log.Printf("telegram bot: webhook set to %s", webhookURL)
logger.Info("telegram bot webhook set", "url", webhookURL)
}
} else {
go func() {
log.Println("telegram bot: starting long-poll loop")
logger.Info("telegram bot starting long-poll loop")
bot.RunLoop(func(u telegram.Update) {
tgH.ProcessUpdate(u)
})
@ -247,8 +253,9 @@ func main() {
}
addr := fmt.Sprintf(":%d", cfg.AppPort)
log.Printf("listening on %s", addr)
logger.Info("listening", "addr", addr)
if err := app.Listen(addr); err != nil {
log.Fatalf("listen: %v", err)
logger.Error("listen failed", "err", err)
os.Exit(1)
}
}

View File

@ -30,9 +30,14 @@ func NewRepo(pool *pgxpool.Pool) *Repo {
func (r *Repo) Log(ctx context.Context, e *Event) error {
meta := []byte("null")
if e.Metadata != nil {
meta = []byte(`{}`) // упрощённо
// NOTE: для prod надо marshal JSON; пропускаем пока
if len(e.Metadata) > 0 {
// Если Metadata уже валидный JSON — пишем как есть
// Иначе (хэштег, текст) — обернём в JSON-объект
meta = e.Metadata
// если начинается не с { или [, оборачиваем
if len(meta) == 0 || (meta[0] != '{' && meta[0] != '[' && meta[0] != '"') {
meta = []byte(`{"data":` + string(meta) + `}`)
}
}
return r.pool.QueryRow(ctx, `
INSERT INTO audit_log (user_id, action, target_type, target_id, metadata, ip, user_agent)

View File

@ -25,6 +25,7 @@ type Message struct {
Body string `json:"body"`
PhotoURL string `json:"photo_url,omitempty"`
Read bool `json:"read"`
ReadAt *time.Time `json:"read_at,omitempty"`
Edited bool `json:"edited"`
Deleted bool `json:"deleted"`
CreatedAt time.Time `json:"created_at"`
@ -39,34 +40,24 @@ func NewRepo(pool *pgxpool.Pool) *Repo {
return &Repo{pool: pool}
}
// EnsureChat — нормализуем user_a < user_b, чтобы избежать дублей
// EnsureChat — race-safe: ON CONFLICT DO NOTHING + RETURNING.
// Нормализуем user_a < user_b, чтобы избежать дублей пары.
func (r *Repo) EnsureChat(ctx context.Context, me, other uuid.UUID) (*Chat, error) {
a, b := me, other
if a.String() > b.String() {
a, b = b, a
}
// ищем существующий
c := &Chat{}
err := r.pool.QueryRow(ctx, `
SELECT id, user_a, user_b, created_at, last_msg_at
FROM chats WHERE user_a=$1 AND user_b=$2`, a, b,
).Scan(&c.ID, &c.UserA, &c.UserB, &c.CreatedAt, &c.LastMsgAt)
if err == nil {
return c, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return nil, err
}
// создаём
err = r.pool.QueryRow(ctx, `
INSERT INTO chats (user_a, user_b) VALUES ($1, $2)
RETURNING id, created_at, last_msg_at`,
ON CONFLICT (user_a, user_b) DO UPDATE
SET last_msg_at = chats.last_msg_at -- noop, чтобы вернулось
RETURNING id, user_a, user_b, created_at, last_msg_at`,
a, b,
).Scan(&c.ID, &c.CreatedAt, &c.LastMsgAt)
).Scan(&c.ID, &c.UserA, &c.UserB, &c.CreatedAt, &c.LastMsgAt)
if err != nil {
return nil, err
}
c.UserA, c.UserB = a, b
return c, nil
}
@ -140,20 +131,34 @@ func (r *Repo) ListChats(ctx context.Context, me uuid.UUID, limit, offset int) (
return out, rows.Err()
}
// Send — атомарно: INSERT message + UPDATE chat.last_msg_at в одной транзакции
func (r *Repo) Send(ctx context.Context, chatID, senderID uuid.UUID, body, photoURL string) (*Message, error) {
m := &Message{}
err := r.pool.QueryRow(ctx, `
INSERT INTO messages (chat_id, sender_id, body, photo_url)
VALUES ($1, $2, NULLIF($3, ''), NULLIF($4, ''))
RETURNING id, chat_id, sender_id, COALESCE(body, ''), COALESCE(photo_url, ''),
read, edited, deleted, created_at, updated_at`,
chatID, senderID, body, photoURL,
).Scan(&m.ID, &m.ChatID, &m.SenderID, &m.Body, &m.PhotoURL,
&m.Read, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt)
tx, err := r.pool.Begin(ctx)
if err != nil {
return nil, err
}
_, _ = r.pool.Exec(ctx, `UPDATE chats SET last_msg_at = NOW() WHERE id=$1`, chatID)
defer tx.Rollback(ctx)
err = tx.QueryRow(ctx, `
INSERT INTO messages (chat_id, sender_id, body, photo_url)
VALUES ($1, $2, NULLIF($3, ''), NULLIF($4, ''))
RETURNING id, chat_id, sender_id, COALESCE(body, ''), COALESCE(photo_url, ''),
read, read_at, edited, deleted, created_at, updated_at`,
chatID, senderID, body, photoURL,
).Scan(&m.ID, &m.ChatID, &m.SenderID, &m.Body, &m.PhotoURL,
&m.Read, &m.ReadAt, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt)
if err != nil {
return nil, err
}
if _, err := tx.Exec(ctx, `UPDATE chats SET last_msg_at = NOW() WHERE id=$1`, chatID); err != nil {
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
return m, nil
}
@ -163,7 +168,7 @@ func (r *Repo) ListMessages(ctx context.Context, chatID uuid.UUID, limit, offset
}
rows, err := r.pool.Query(ctx, `
SELECT id, chat_id, sender_id, COALESCE(body, ''), COALESCE(photo_url, ''),
read, edited, deleted, created_at, updated_at
read, read_at, edited, deleted, created_at, updated_at
FROM messages
WHERE chat_id=$1
ORDER BY created_at ASC
@ -176,7 +181,7 @@ func (r *Repo) ListMessages(ctx context.Context, chatID uuid.UUID, limit, offset
for rows.Next() {
var m Message
if err := rows.Scan(&m.ID, &m.ChatID, &m.SenderID, &m.Body, &m.PhotoURL,
&m.Read, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt); err != nil {
&m.Read, &m.ReadAt, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt); err != nil {
return nil, err
}
out = append(out, m)
@ -203,10 +208,10 @@ func (r *Repo) GetMessage(ctx context.Context, msgID uuid.UUID) (*Message, error
m := &Message{}
err := r.pool.QueryRow(ctx, `
SELECT id, chat_id, sender_id, COALESCE(body, ''), COALESCE(photo_url, ''),
read, edited, deleted, created_at, updated_at
read, read_at, edited, deleted, created_at, updated_at
FROM messages WHERE id=$1`, msgID,
).Scan(&m.ID, &m.ChatID, &m.SenderID, &m.Body, &m.PhotoURL,
&m.Read, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt)
&m.Read, &m.ReadAt, &m.Edited, &m.Deleted, &m.CreatedAt, &m.UpdatedAt)
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
@ -231,13 +236,34 @@ func (r *Repo) DeleteMessage(ctx context.Context, msgID, senderID uuid.UUID) err
return nil
}
// MarkRead — атомарно с проверкой членства. Возвращает chat_id для WS broadcast.
func (r *Repo) MarkRead(ctx context.Context, chatID, meID uuid.UUID) error {
_, err := r.pool.Exec(ctx, `
UPDATE messages SET read=TRUE
WHERE chat_id=$1 AND sender_id != $2 AND read=FALSE`,
// member-check + update в одном UPDATE с подзапросом
res, err := r.pool.Exec(ctx, `
UPDATE messages SET read=TRUE, read_at=NOW()
WHERE chat_id=$1 AND sender_id != $2 AND read=FALSE
AND EXISTS (SELECT 1 FROM chats c
WHERE c.id=$1 AND (c.user_a=$2 OR c.user_b=$2))`,
chatID, meID,
)
if err != nil {
return err
}
if res.RowsAffected() == 0 {
// не ошибка, просто либо нет сообщений, либо юзер не в чате
// проверим явно membership для возврата ошибки
var isMember bool
if err := r.pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM chats WHERE id=$1 AND (user_a=$2 OR user_b=$2))`,
chatID, meID,
).Scan(&isMember); err != nil {
return err
}
if !isMember {
return errors.New("not a member of this chat")
}
}
return nil
}
// Block / Unblock

View File

@ -147,6 +147,9 @@ func (h *ChatHandlers) SendMessage(c *fiber.Ctx) error {
"sender_id": m.SenderID,
"body": m.Body,
"photo_url": m.PhotoURL,
"read": m.Read,
"read_at": m.ReadAt,
"edited": m.Edited,
"created_at": m.CreatedAt,
},
Time: time.Now(),
@ -160,6 +163,9 @@ func (h *ChatHandlers) SendMessage(c *fiber.Ctx) error {
"sender_id": m.SenderID,
"body": m.Body,
"photo_url": m.PhotoURL,
"read": m.Read,
"read_at": m.ReadAt,
"edited": m.Edited,
"created_at": m.CreatedAt,
},
Time: time.Now(),
@ -278,7 +284,24 @@ func (h *ChatHandlers) MarkRead(c *fiber.Ctx) error {
return c.Status(fiber.StatusBadRequest).JSON(fiber.Map{"error": "bad chat id"})
}
if err := h.Chat.MarkRead(c.UserContext(), chatID, me); err != nil {
return c.Status(fiber.StatusInternalServerError).JSON(fiber.Map{"error": "db error"})
return c.Status(fiber.StatusForbidden).JSON(fiber.Map{"error": err.Error()})
}
// оповестим отправителя через WS что сообщения прочитаны
chat, _ := h.Chat.GetByID(c.UserContext(), chatID)
if chat != nil {
other := chat.UserA
if other == me {
other = chat.UserB
}
h.Hub.SendTo(other, ws.Outgoing{
Type: "message_read",
Payload: fiber.Map{
"chat_id": chatID,
"reader": me,
"read_at": time.Now(),
},
Time: time.Now(),
})
}
return c.JSON(fiber.Map{"ok": true})
}

View File

@ -76,6 +76,8 @@ func (r *Repo) Get(ctx context.Context, userID uuid.UUID) (*Location, error) {
}
// Nearby — пользователи в радиусе radius км (по haversine), visible, не старше 1 часа
// Оптимизация: фильтруем сначала по bounding-box (lat/lng), потом haversine только для отфильтрованных.
// Это устраняет двойной расчёт haversine и даёт индекс если добавить GiST/btree на lat,lng.
type Nearby struct {
UserID uuid.UUID
Name string
@ -89,13 +91,18 @@ func (r *Repo) Nearby(ctx context.Context, myID uuid.UUID, lat, lng, radiusKm fl
if limit <= 0 || limit > 500 {
limit = 100
}
// bounding box: 1° lat ≈ 111 км, 1° lng ≈ 111 км * cos(lat)
latDelta := radiusKm / 111.0
lngDelta := radiusKm / (111.0 * math.Cos(lat*math.Pi/180))
if latDelta < 0.001 {
latDelta = 0.001
}
if lngDelta < 0.001 {
lngDelta = 0.001
}
rows, err := r.pool.Query(ctx, `
SELECT u.id, u.name, u.photo_url, ul.lat, ul.lng,
2 * 6371 * asin(sqrt(
sin(radians((ul.lat - $1) / 2))^2 +
cos(radians($1)) * cos(radians(ul.lat)) *
sin(radians((ul.lng - $2) / 2))^2
)) AS dist_km
WITH bbox AS (
SELECT u.id, u.name, u.photo_url, ul.lat, ul.lng
FROM user_locations ul
JOIN users u ON u.id = ul.user_id
WHERE ul.visible = TRUE
@ -103,14 +110,24 @@ func (r *Repo) Nearby(ctx context.Context, myID uuid.UUID, lat, lng, radiusKm fl
AND u.id != $3
AND ul.lat IS NOT NULL AND ul.lng IS NOT NULL
AND ul.updated_at > NOW() - INTERVAL '1 hour'
AND 2 * 6371 * asin(sqrt(
sin(radians((ul.lat - $1) / 2))^2 +
cos(radians($1)) * cos(radians(ul.lat)) *
sin(radians((ul.lng - $2) / 2))^2
AND ul.lat BETWEEN $1 - $6 AND $1 + $6
AND ul.lng BETWEEN $2 - $7 AND $2 + $7
)
SELECT id, name, photo_url, lat, lng,
2 * 6371 * asin(sqrt(
sin(radians((lat - $1) / 2))^2 +
cos(radians($1)) * cos(radians(lat)) *
sin(radians((lng - $2) / 2))^2
)) * 1000 AS dist_m
FROM bbox
WHERE 2 * 6371 * asin(sqrt(
sin(radians((lat - $1) / 2))^2 +
cos(radians($1)) * cos(radians(lat)) *
sin(radians((lng - $2) / 2))^2
)) <= $4
ORDER BY dist_km
ORDER BY dist_m
LIMIT $5`,
lat, lng, myID, radiusKm, limit,
lat, lng, myID, radiusKm, limit, latDelta, lngDelta,
)
if err != nil {
return nil, err
@ -127,7 +144,6 @@ func (r *Repo) Nearby(ctx context.Context, myID uuid.UUID, lat, lng, radiusKm fl
if photo != nil {
n.PhotoURL = *photo
}
n.DistanceM = n.DistanceM * 1000
out = append(out, n)
}
return out, rows.Err()

View File

@ -0,0 +1 @@
ALTER TABLE messages DROP COLUMN IF EXISTS read_at;

View File

@ -0,0 +1,3 @@
-- read_at timestamp для галочек "прочитано"
ALTER TABLE messages ADD COLUMN IF NOT EXISTS read_at TIMESTAMPTZ;
UPDATE messages SET read_at = updated_at WHERE read = TRUE AND read_at IS NULL;