From 76a53a715db1457e1133ff7ae6e46a8c88f75cdc Mon Sep 17 00:00:00 2001 From: Aleksandr Soloshenko Date: Sat, 25 Jul 2026 08:13:43 +0700 Subject: [PATCH 1/2] [inbox] add module --- .golangci.yml | 1 + README.md | 4 +- api/mobile.http | 94 +++++++++++ go.mod | 2 +- go.sum | 4 +- internal/config/module.go | 4 + .../sms-gateway/handlers/inbox/3rdparty.go | 97 +++++++++++- internal/sms-gateway/handlers/inbox/errors.go | 26 +++ internal/sms-gateway/handlers/inbox/mobile.go | 93 +++++++++++ internal/sms-gateway/handlers/inbox/params.go | 27 ++++ .../sms-gateway/handlers/inbox/permissions.go | 1 + .../sms-gateway/handlers/messages/params.go | 13 +- internal/sms-gateway/handlers/mobile.go | 5 + internal/sms-gateway/handlers/module.go | 1 + internal/sms-gateway/inbox/config.go | 4 + internal/sms-gateway/inbox/domain.go | 100 ++++++++++++ internal/sms-gateway/inbox/errors.go | 9 ++ internal/sms-gateway/inbox/models.go | 122 ++++++++++++++ internal/sms-gateway/inbox/module.go | 12 +- internal/sms-gateway/inbox/repository.go | 149 ++++++++++++++++++ internal/sms-gateway/inbox/service.go | 59 ++++++- .../20260724000000_create_inbox_tables.sql | 44 ++++++ internal/sms-gateway/openapi/docs.go | 126 ++++++++++++++- internal/worker/config/config.go | 10 ++ internal/worker/config/module.go | 9 ++ internal/worker/tasks/inbox/cleanup.go | 56 +++++++ internal/worker/tasks/inbox/config.go | 12 ++ internal/worker/tasks/inbox/module.go | 22 +++ internal/worker/tasks/module.go | 2 + 29 files changed, 1076 insertions(+), 32 deletions(-) create mode 100644 internal/sms-gateway/handlers/inbox/errors.go create mode 100644 internal/sms-gateway/handlers/inbox/mobile.go create mode 100644 internal/sms-gateway/handlers/inbox/params.go create mode 100644 internal/sms-gateway/inbox/config.go create mode 100644 internal/sms-gateway/inbox/domain.go create mode 100644 internal/sms-gateway/inbox/errors.go create mode 100644 internal/sms-gateway/inbox/models.go create mode 100644 internal/sms-gateway/inbox/repository.go create mode 100644 internal/sms-gateway/models/migrations/mysql/20260724000000_create_inbox_tables.sql create mode 100644 internal/worker/tasks/inbox/cleanup.go create mode 100644 internal/worker/tasks/inbox/config.go create mode 100644 internal/worker/tasks/inbox/module.go diff --git a/.golangci.yml b/.golangci.yml index 1c6f45c8..bc7ef2f4 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -462,6 +462,7 @@ linters: wrapcheck: extra-ignore-sigs: - .JSON( + - .Send( - .SendStatus( exclusions: diff --git a/README.md b/README.md index 1cb46d82..521b3bd1 100644 --- a/README.md +++ b/README.md @@ -165,8 +165,8 @@ The following scopes are available for token generation: - `devices:delete` - Delete devices - `devices:list` - List connected devices -- `inbox:list` - List incoming messages with filters -- `inbox:read` - Read incoming messages +- `inbox:list` - List inbox messages with filters +- `inbox:read` - Read inbox messages - `logs:read` - Read server logs - `messages:export` - Export messages - `messages:list` - List messages diff --git a/api/mobile.http b/api/mobile.http index 519e8483..d403de5a 100644 --- a/api/mobile.http +++ b/api/mobile.http @@ -81,3 +81,97 @@ Authorization: Bearer {{mobileToken}} ### GET {{baseUrl}}/events HTTP/1.1 Authorization: Bearer {{mobileToken}} + +### +POST {{baseUrl}}/inbox HTTP/1.1 +Authorization: Bearer {{mobileToken}} +Content-Type: application/json + +[ + { + "id": "sms-001", + "type": "SMS", + "sender": "+79990001234", + "recipient": "+79990005678", + "simNumber": 1, + "content": "a2V5X2VuY3J5cHRlZF9jb250ZW50", + "isEncrypted": true, + "createdAt": "2026-01-15T10:30:00Z" + } +] + +### +POST {{baseUrl}}/inbox HTTP/1.1 +Authorization: Bearer {{mobileToken}} +Content-Type: application/json + +[ + { + "id": "mms-001", + "type": "MMS", + "sender": "+79990001234", + "recipient": "+79990005678", + "simNumber": 1, + "content": "a2V5X2VuY3J5cHRlZF9jb250ZW50", + "isEncrypted": true, + "createdAt": "2026-01-15T11:00:00Z", + "attachments": [ + { + "partId": 1, + "contentType": "image/jpeg", + "name": "enc_photo.jpg", + "size": 102400, + "data": "base64encodeddata" + }, + { + "partId": 2, + "contentType": "text/plain", + "name": "enc_text.txt", + "data": "base64encodeddata" + } + ] + } +] + +### +POST {{baseUrl}}/inbox HTTP/1.1 +Authorization: Bearer {{mobileToken}} +Content-Type: application/json + +[ + { + "id": "sms-batch-001", + "type": "SMS", + "sender": "+79990001111", + "content": "a2V5X2VuY3J5cHRlZF9jb250ZW50XzE=", + "isEncrypted": true, + "createdAt": "2026-01-15T12:00:00Z" + }, + { + "id": "sms-batch-002", + "type": "DATA_SMS", + "sender": "+79990002222", + "recipient": "+79990005678", + "simNumber": 2, + "content": "a2V5X2VuY3J5cHRlZF9jb250ZW50XzI=", + "isEncrypted": true, + "createdAt": "2026-01-15T12:01:00Z" + }, + { + "id": "mms-batch-001", + "type": "MMS_DOWNLOADED", + "sender": "+79990003333", + "content": "a2V5X2VuY3J5cHRlZF9jb250ZW50XzM=", + "isEncrypted": true, + "createdAt": "2026-01-15T12:02:00Z", + "attachments": [ + { + "partId": 1, + "contentType": "image/png", + "name": "enc_image.png", + "size": 204800, + "data": "base64encodeddata" + } + ] + } +] diff --git a/go.mod b/go.mod index 7d5b64cd..15356de3 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.25.8 require ( firebase.google.com/go/v4 v4.21.0 - github.com/android-sms-gateway/client-go v1.14.5 + github.com/android-sms-gateway/client-go v1.14.6-0.20260818003959-ca8711dc4282 github.com/ansrivas/fiberprometheus/v2 v2.17.0 github.com/capcom6/go-helpers v0.4.0 github.com/capcom6/go-infra-fx v0.5.9 diff --git a/go.sum b/go.sum index c1dc690c..c08f9f2e 100644 --- a/go.sum +++ b/go.sum @@ -38,8 +38,8 @@ github.com/KyleBanks/depth v1.2.1 h1:5h8fQADFrWtarTdtDudMmGsC7GPbOAu6RVB3ffsVFHc github.com/KyleBanks/depth v1.2.1/go.mod h1:jzSb9d0L43HxTQfT+oSA1EEp2q+ne2uh6XgeJcm8brE= github.com/MicahParks/keyfunc v1.9.0 h1:lhKd5xrFHLNOWrDc4Tyb/Q1AJ4LCzQ48GVJyVIID3+o= github.com/MicahParks/keyfunc v1.9.0/go.mod h1:IdnCilugA0O/99dW+/MkvlyrsX8+L8+x95xuVNtM5jw= -github.com/android-sms-gateway/client-go v1.14.5 h1:CtyXAHPdyDtCFTzeVPt5EUfnWhQ2RsIWbCUk52tbGjw= -github.com/android-sms-gateway/client-go v1.14.5/go.mod h1:DQsReciU1xcaVW3T5Z2bqslNdsAwCFCtghawmA6g6L4= +github.com/android-sms-gateway/client-go v1.14.6-0.20260818003959-ca8711dc4282 h1:L02pWbWnYwZJAoxdtyKi9tru7HDCLAGyJcFInwws8TQ= +github.com/android-sms-gateway/client-go v1.14.6-0.20260818003959-ca8711dc4282/go.mod h1:DQsReciU1xcaVW3T5Z2bqslNdsAwCFCtghawmA6g6L4= github.com/andybalholm/brotli v1.2.2 h1:HzTuoo2ErYQqf5qvcJInB8uvqSVxRttzkFexPWtnceM= github.com/andybalholm/brotli v1.2.2/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY= github.com/ansrivas/fiberprometheus/v2 v2.17.0 h1:p0gqs5LsSCWGoSFF44fCJkyU+XcE6TLRqEMu80b2iCo= diff --git a/internal/config/module.go b/internal/config/module.go index 21ba3729..a7cb325e 100644 --- a/internal/config/module.go +++ b/internal/config/module.go @@ -5,6 +5,7 @@ import ( "time" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers" + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" "github.com/android-sms-gateway/server/internal/sms-gateway/jwt" "github.com/android-sms-gateway/server/internal/sms-gateway/modules/auth" "github.com/android-sms-gateway/server/internal/sms-gateway/modules/devices" @@ -128,6 +129,9 @@ func Module() fx.Option { fx.Provide(func(_ Config) devices.Config { return devices.Config{} }), + fx.Provide(func(_ Config) inbox.Config { + return inbox.Config{} + }), fx.Provide(func(cfg Config) sse.Config { return sse.NewConfig( sse.WithKeepAlivePeriod(time.Duration(cfg.SSE.KeepAlivePeriodSeconds) * time.Second), diff --git a/internal/sms-gateway/handlers/inbox/3rdparty.go b/internal/sms-gateway/handlers/inbox/3rdparty.go index bb430481..e49c8c52 100644 --- a/internal/sms-gateway/handlers/inbox/3rdparty.go +++ b/internal/sms-gateway/handlers/inbox/3rdparty.go @@ -1,6 +1,10 @@ package inbox import ( + "errors" + "fmt" + "strconv" + "github.com/android-sms-gateway/client-go/smsgateway" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/base" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/middlewares/permissions" @@ -8,6 +12,7 @@ import ( "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" "github.com/go-playground/validator/v10" "github.com/gofiber/fiber/v2" + "github.com/samber/lo" "go.uber.org/zap" ) @@ -33,23 +38,25 @@ func NewThirdPartyController( } func (h *ThirdPartyController) Register(router fiber.Router) { + router.Use(errorHandler) router.Get("", permissions.RequireScope(ScopeList), userauth.WithUserID(h.list)) + router.Get("/:id/attachments/:partId", permissions.RequireScope(ScopeRead), userauth.WithUserID(h.getAttachment)) router.Post("/refresh", permissions.RequireScope(ScopeRefresh), userauth.WithUserID(h.refresh)) } -// @Summary Get incoming messages -// @Description Retrieves incoming messages with filtering and pagination. +// @Summary Get inbox messages +// @Description Retrieves inbox messages with filtering and pagination. // @Security ApiAuth // @Security JWTAuth // @Tags User, Inbox // @Produce json -// @Param type query string false "Filter incoming messages by type" Enums(SMS,DATA_SMS,MMS,MMS_DOWNLOADED) +// @Param type query string false "Filter inbox messages by type" Enums(SMS,DATA_SMS,MMS,MMS_DOWNLOADED) // @Param limit query int false "Maximum number of messages to return" minimum(1) maximum(500) default(50) // @Param offset query int false "Number of messages to skip" minimum(0) default(0) // @Param from query string false "Start of date range (ISO 8601)" Format(date-time) // @Param to query string false "End of date range (ISO 8601)" Format(date-time) // @Param deviceId query string false "Device ID" -// @Success 200 {array} smsgateway.IncomingMessage "A list of incoming messages" +// @Success 200 {array} smsgateway.IncomingMessage "A list of inbox messages" // @Header 200 {integer} X-Total-Count "Total number of items available" // @Failure 400 {object} smsgateway.ErrorResponse "Invalid request" // @Failure 401 {object} smsgateway.ErrorResponse "Unauthorized" @@ -58,9 +65,46 @@ func (h *ThirdPartyController) Register(router fiber.Router) { // @Failure 501 {object} smsgateway.ErrorResponse "Not implemented" // @Router /3rdparty/v1/inbox [get] // -// Get incoming messages. -func (h *ThirdPartyController) list(_ string, _ *fiber.Ctx) error { - return fiber.NewError(fiber.StatusNotImplemented, "Inbox API is not implemented yet") +// Get inbox messages. +func (h *ThirdPartyController) list(userID string, c *fiber.Ctx) error { + var params thirdPartyListParams + if err := c.QueryParser(¶ms); err != nil { + h.Logger.Error("failed to parse query parameters", zap.Error(err)) + return fiber.NewError(fiber.StatusBadRequest, "failed to parse query parameters") + } + + messages, total, err := h.inboxSvc.List(userID, params.toFilter(), params.toOptions()) + if err != nil { + h.Logger.Error("failed to list inbox messages", zap.Error(err), zap.String("user_id", userID)) + return fiber.NewError(fiber.StatusInternalServerError, "failed to list inbox messages") + } + + result := make([]smsgateway.IncomingMessage, len(messages)) + for i, m := range messages { + atts := make([]smsgateway.InboxAttachment, len(m.Attachments)) + for j, a := range m.Attachments { + atts[j] = smsgateway.InboxAttachment{ + PartID: a.PartID, + Name: a.Name, + Size: lo.FromPtrOr(a.Size, 0), + ContentType: a.ContentType, + } + } + result[i] = smsgateway.IncomingMessage{ + ID: m.ExtID, + Type: smsgateway.IncomingMessageType(m.Type), + Sender: m.Sender, + Recipient: m.Recipient, + SimNumber: m.SimNumber, + ContentPreview: m.Content, + IsEncrypted: m.IsEncrypted, + CreatedAt: m.CreatedAt, + Attachments: atts, + } + } + + c.Set("X-Total-Count", strconv.Itoa(int(total))) + return c.JSON(result) } // @Summary Request inbox messages refresh @@ -98,3 +142,42 @@ func (h *ThirdPartyController) refresh(userID string, c *fiber.Ctx) error { return c.SendStatus(fiber.StatusAccepted) } + +// @Summary Get attachment +// @Description Downloads an attachment from an inbox message by message ID and part ID. +// @Security ApiAuth +// @Security JWTAuth +// @Tags User, Inbox +// @Produce application/octet-stream +// @Param id path string true "Inbox message ID" +// @Param partId path int true "Attachment part ID" +// @Success 200 {file} binary "Attachment file" +// @Failure 400 {object} smsgateway.ErrorResponse "Invalid request" +// @Failure 401 {object} smsgateway.ErrorResponse "Unauthorized" +// @Failure 403 {object} smsgateway.ErrorResponse "Forbidden" +// @Failure 404 {object} smsgateway.ErrorResponse "Not found" +// @Failure 500 {object} smsgateway.ErrorResponse "Internal server error" +// @Router /3rdparty/v1/inbox/{id}/attachments/{partId} [get] +// +// Get attachment. +func (h *ThirdPartyController) getAttachment(userID string, c *fiber.Ctx) error { + id := c.Params("id") + partID, err := strconv.ParseInt(c.Params("partId"), 10, 64) + if err != nil { + return fiber.NewError(fiber.StatusBadRequest, "invalid partId") + } + + att, err := h.inboxSvc.GetAttachment(c.Context(), userID, id, partID) + if err != nil { + if errors.Is(err, inbox.ErrNotFound) { + return fiber.NewError(fiber.StatusNotFound, "attachment not found") + } + h.Logger.Error("failed to get attachment", zap.Error(err), zap.String("user_id", userID)) + return fiber.NewError(fiber.StatusInternalServerError, "failed to get attachment") + } + + c.Set("Content-Disposition", fmt.Sprintf(`attachment; filename="%s"`, att.Name)) + c.Set("Content-Type", att.ContentType) + + return c.Send(att.Data) +} diff --git a/internal/sms-gateway/handlers/inbox/errors.go b/internal/sms-gateway/handlers/inbox/errors.go new file mode 100644 index 00000000..c20d9c35 --- /dev/null +++ b/internal/sms-gateway/handlers/inbox/errors.go @@ -0,0 +1,26 @@ +package inbox + +import ( + "errors" + + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" + "github.com/gofiber/fiber/v2" +) + +func errorHandler(c *fiber.Ctx) error { + err := c.Next() + if err == nil { + return nil + } + + switch { + case errors.Is(err, inbox.ErrNotEncrypted): + return fiber.NewError(fiber.StatusBadRequest, err.Error()) + case errors.Is(err, inbox.ErrNotFound): + return fiber.NewError(fiber.StatusNotFound, err.Error()) + case errors.Is(err, inbox.ErrEmptyBatch): + return fiber.NewError(fiber.StatusBadRequest, err.Error()) + } + + return err //nolint:wrapcheck // passed through to fiber's error handler +} diff --git a/internal/sms-gateway/handlers/inbox/mobile.go b/internal/sms-gateway/handlers/inbox/mobile.go new file mode 100644 index 00000000..ac55300e --- /dev/null +++ b/internal/sms-gateway/handlers/inbox/mobile.go @@ -0,0 +1,93 @@ +package inbox + +import ( + "fmt" + + "github.com/android-sms-gateway/client-go/smsgateway" + "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/base" + "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/middlewares/deviceauth" + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" + "github.com/android-sms-gateway/server/internal/sms-gateway/modules/devices" + "github.com/go-playground/validator/v10" + "github.com/gofiber/fiber/v2" + "go.uber.org/zap" +) + +type MobileController struct { + base.Handler + + inboxSvc *inbox.Service +} + +func NewMobileController( + inboxSvc *inbox.Service, + logger *zap.Logger, + validator *validator.Validate, +) *MobileController { + return &MobileController{ + Handler: base.Handler{ + Logger: logger, + Validator: validator, + }, + inboxSvc: inboxSvc, + } +} + +func (h *MobileController) Register(router fiber.Router) { + router.Use(errorHandler) + router.Post("", deviceauth.WithDevice(h.post)) +} + +// @Summary Upload inbox messages +// @Description Stores a batch of encrypted inbox messages from the device +// @Security MobileToken +// @Tags Device, Inbox +// @Accept json +// @Param request body smsgateway.MobilePostInboxRequest true "Batch of inbox messages" +// @Success 201 "Created" +// @Failure 400 {object} smsgateway.ErrorResponse "Invalid request" +// @Failure 500 {object} smsgateway.ErrorResponse "Internal server error" +// @Router /mobile/v1/inbox [post] +// +// Upload inbox messages. +func (h *MobileController) post(device devices.Device, c *fiber.Ctx) error { + req := make(smsgateway.MobilePostInboxRequest, 0) + if err := h.BodyParserValidator(c, &req); err != nil { + return fiber.NewError(fiber.StatusBadRequest, err.Error()) + } + + msgs := make([]inbox.MessageInput, 0, len(req)) + for _, r := range req { + atts := make([]inbox.AttachmentInput, 0, len(r.Attachments)) + for _, a := range r.Attachments { + atts = append(atts, inbox.AttachmentInput{ + PartID: a.PartID, + ContentType: a.ContentType, + Name: a.Name, + Size: a.Size, + Data: a.Data, + }) + } + + msgs = append(msgs, inbox.MessageInput{ + MessageBody: inbox.MessageBody{ + ExtID: r.ID, + Type: inbox.MessageType(r.Type), + Sender: r.Sender, + Recipient: r.Recipient, + SimNumber: r.SimNumber, + Content: r.Content, + IsEncrypted: r.IsEncrypted, + CreatedAt: r.CreatedAt, + }, + + Attachments: atts, + }) + } + + if err := h.inboxSvc.InsertBatch(c.Context(), device.ID, msgs); err != nil { + return fmt.Errorf("failed to insert inbox messages: %w", err) + } + + return c.SendStatus(fiber.StatusCreated) +} diff --git a/internal/sms-gateway/handlers/inbox/params.go b/internal/sms-gateway/handlers/inbox/params.go new file mode 100644 index 00000000..aea1f35f --- /dev/null +++ b/internal/sms-gateway/handlers/inbox/params.go @@ -0,0 +1,27 @@ +package inbox + +import ( + "github.com/android-sms-gateway/client-go/smsgateway" + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" + "github.com/samber/lo" +) + +type thirdPartyListParams smsgateway.ListInboxOptions + +func (p thirdPartyListParams) toFilter() inbox.ListFilter { + return inbox.ListFilter{ + DeviceID: lo.FromPtr(p.DeviceID), + Type: inbox.MessageType(lo.FromPtr(p.Type)), + StartDate: lo.FromPtr(p.From), + EndDate: lo.FromPtr(p.To), + } +} + +func (p thirdPartyListParams) toOptions() inbox.ListOptions { + const defaultLimit = 50 + + return inbox.ListOptions{ + Limit: lo.FromPtrOr(p.Limit, defaultLimit), + Offset: lo.FromPtr(p.Offset), + } +} diff --git a/internal/sms-gateway/handlers/inbox/permissions.go b/internal/sms-gateway/handlers/inbox/permissions.go index 3c98461c..189ffd24 100644 --- a/internal/sms-gateway/handlers/inbox/permissions.go +++ b/internal/sms-gateway/handlers/inbox/permissions.go @@ -5,4 +5,5 @@ import "github.com/android-sms-gateway/client-go/smsgateway" const ( ScopeList = smsgateway.ScopeInboxList ScopeRefresh = smsgateway.ScopeInboxRefresh + ScopeRead = smsgateway.ScopeInboxRead ) diff --git a/internal/sms-gateway/handlers/messages/params.go b/internal/sms-gateway/handlers/messages/params.go index 8afc0051..77b31787 100644 --- a/internal/sms-gateway/handlers/messages/params.go +++ b/internal/sms-gateway/handlers/messages/params.go @@ -3,6 +3,7 @@ package messages import ( "github.com/android-sms-gateway/client-go/smsgateway" "github.com/android-sms-gateway/server/internal/sms-gateway/modules/messages" + "github.com/samber/lo" ) // thirdPartyPostQueryParams aliases smsgateway.SendOptions so that the query @@ -36,7 +37,7 @@ func (p *thirdPartyGetQueryParams) ToFilter() messages.SelectFilter { } func (p *thirdPartyGetQueryParams) ToOptions() messages.SelectOptions { - const maxLimit = 100 + const defaultLimit = 50 var options messages.SelectOptions options.WithRecipients = true @@ -46,15 +47,7 @@ func (p *thirdPartyGetQueryParams) ToOptions() messages.SelectOptions { options.WithContent = *p.IncludeContent } - if p.Limit != nil { - options.Limit = max(min(*p.Limit, maxLimit), 1) - } else { - options.Limit = 50 - } - - if p.Offset != nil { - options.Offset = max(*p.Offset, 0) - } + options.Limit = lo.FromPtrOr(p.Limit, int(defaultLimit)) if p.Sort != nil { switch *p.Sort { diff --git a/internal/sms-gateway/handlers/mobile.go b/internal/sms-gateway/handlers/mobile.go index dbda34b2..8c6cf11d 100644 --- a/internal/sms-gateway/handlers/mobile.go +++ b/internal/sms-gateway/handlers/mobile.go @@ -9,6 +9,7 @@ import ( "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/base" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/converters" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/events" + "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/inbox" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/messages" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/middlewares/deviceauth" "github.com/android-sms-gateway/server/internal/sms-gateway/handlers/middlewares/userauth" @@ -37,6 +38,7 @@ type mobileHandler struct { webhooksCtrl *webhooks.MobileController settingsCtrl *settings.MobileController eventsCtrl *events.MobileController + inboxCtrl *inbox.MobileController idGen func() string } @@ -50,6 +52,7 @@ func newMobileHandler( webhooksCtrl *webhooks.MobileController, settingsCtrl *settings.MobileController, eventsCtrl *events.MobileController, + inboxCtrl *inbox.MobileController, logger *zap.Logger, validator *validator.Validate, @@ -70,6 +73,7 @@ func newMobileHandler( webhooksCtrl: webhooksCtrl, settingsCtrl: settingsCtrl, eventsCtrl: eventsCtrl, + inboxCtrl: inboxCtrl, idGen: idGen, } @@ -124,6 +128,7 @@ func (h *mobileHandler) Register(router fiber.Router) { h.webhooksCtrl.Register(router.Group("/webhooks")) h.settingsCtrl.Register(router.Group("/settings")) h.eventsCtrl.Register(router.Group("/events")) + h.inboxCtrl.Register(router.Group("/inbox")) } // @Summary Get device information diff --git a/internal/sms-gateway/handlers/module.go b/internal/sms-gateway/handlers/module.go index 66a1a6df..17696c52 100644 --- a/internal/sms-gateway/handlers/module.go +++ b/internal/sms-gateway/handlers/module.go @@ -37,6 +37,7 @@ func Module() fx.Option { settings.NewThirdPartyController, settings.NewMobileController, inbox.NewThirdPartyController, + inbox.NewMobileController, logs.NewThirdPartyController, events.NewMobileController, fx.Private, diff --git a/internal/sms-gateway/inbox/config.go b/internal/sms-gateway/inbox/config.go new file mode 100644 index 00000000..fe21bed5 --- /dev/null +++ b/internal/sms-gateway/inbox/config.go @@ -0,0 +1,4 @@ +package inbox + +type Config struct { +} diff --git a/internal/sms-gateway/inbox/domain.go b/internal/sms-gateway/inbox/domain.go new file mode 100644 index 00000000..4900480b --- /dev/null +++ b/internal/sms-gateway/inbox/domain.go @@ -0,0 +1,100 @@ +package inbox + +import ( + "time" + + "gorm.io/gorm" +) + +// MessageType represents the type of inbox message. +type MessageType string + +const ( + MessageTypeSMS MessageType = "SMS" + MessageTypeDataSMS MessageType = "DATA_SMS" + MessageTypeMMS MessageType = "MMS" + MessageTypeMmsDownloaded MessageType = "MMS_DOWNLOADED" +) + +type MessageBody struct { + ExtID string + Type MessageType + Sender string + Recipient *string + SimNumber *uint8 + Content string + IsEncrypted bool + CreatedAt time.Time +} + +// MessageInput holds the data needed to create an inbox message. +type MessageInput struct { + MessageBody + + Attachments []AttachmentInput +} + +// AttachmentInput holds the data needed to create an inbox attachment. +type AttachmentInput struct { + PartID int64 + ContentType string + Name string + Size *int64 + Data []byte +} + +// Message represents a stored inbox message. +type Message struct { + MessageBody + + ID uint64 + DeviceID string + Attachments []Attachment +} + +// Attachment represents a stored inbox attachment. +type Attachment struct { + AttachmentInput + + ID uint64 +} + +// ListFilter defines filters for listing inbox messages. +type ListFilter struct { + DeviceID string + Type MessageType + StartDate time.Time + EndDate time.Time +} + +func (f ListFilter) apply(query *gorm.DB) *gorm.DB { + if f.DeviceID != "" { + query = query.Where("device_id = ?", f.DeviceID) + } + if f.Type != "" { + query = query.Where("type = ?", f.Type) + } + if !f.StartDate.IsZero() { + query = query.Where("created_at >= ?", f.StartDate) + } + if !f.EndDate.IsZero() { + query = query.Where("created_at <= ?", f.EndDate) + } + return query +} + +// ListOptions defines pagination options for listing inbox messages. +type ListOptions struct { + Limit int + Offset int +} + +func (o ListOptions) apply(query *gorm.DB) *gorm.DB { + if o.Limit > 0 { + query = query.Limit(o.Limit) + } + if o.Offset > 0 { + query = query.Offset(o.Offset) + } + return query +} diff --git a/internal/sms-gateway/inbox/errors.go b/internal/sms-gateway/inbox/errors.go new file mode 100644 index 00000000..5290bd30 --- /dev/null +++ b/internal/sms-gateway/inbox/errors.go @@ -0,0 +1,9 @@ +package inbox + +import "errors" + +var ( + ErrNotEncrypted = errors.New("inbox messages must be encrypted") + ErrNotFound = errors.New("inbox message not found") + ErrEmptyBatch = errors.New("batch must contain at least one message") +) diff --git a/internal/sms-gateway/inbox/models.go b/internal/sms-gateway/inbox/models.go new file mode 100644 index 00000000..f617bc71 --- /dev/null +++ b/internal/sms-gateway/inbox/models.go @@ -0,0 +1,122 @@ +package inbox + +import ( + "fmt" + + "github.com/android-sms-gateway/server/internal/sms-gateway/models" + "github.com/samber/lo" + "gorm.io/gorm" +) + +type messageModel struct { + models.TimedModel + + ID uint64 `gorm:"primaryKey;type:BIGINT UNSIGNED;autoIncrement"` + ExtID string `gorm:"not null;type:varchar(36);uniqueIndex:unq_inbox_ext_device,priority:1"` + DeviceID string `gorm:"not null;type:char(21);uniqueIndex:unq_inbox_ext_device,priority:2;index:idx_inbox_device_created"` + Type MessageType `gorm:"not null;type:enum('SMS','DATA_SMS','MMS','MMS_DOWNLOADED')"` + Sender string `gorm:"not null;type:varchar(512)"` + Recipient *string `gorm:"type:varchar(512)"` + SimNumber *uint8 `gorm:"type:tinyint(1) unsigned"` + Content string `gorm:"not null;type:text"` + IsEncrypted bool `gorm:"not null;type:tinyint(1) unsigned;default:1"` + + Attachments []attachmentModel `gorm:"foreignKey:MessageID;constraint:OnDelete:CASCADE"` +} + +func newMessageModel(deviceID string, input MessageInput) *messageModel { + return &messageModel{ + TimedModel: models.TimedModel{ + CreatedAt: input.CreatedAt, + UpdatedAt: input.CreatedAt, + }, + + ID: 0, + ExtID: input.ExtID, + DeviceID: deviceID, + Type: input.Type, + Sender: input.Sender, + Recipient: input.Recipient, + SimNumber: input.SimNumber, + Content: input.Content, + IsEncrypted: input.IsEncrypted, + + Attachments: lo.Map( + input.Attachments, + func(item AttachmentInput, _ int) attachmentModel { return newAttachmentModel(item) }, + ), + } +} + +func newAttachmentModel(input AttachmentInput) attachmentModel { + return attachmentModel{ + ID: 0, + MessageID: 0, + PartID: input.PartID, + ContentType: input.ContentType, + Name: input.Name, + Size: input.Size, + Data: input.Data, + } +} + +func (*messageModel) TableName() string { + return "inbox" +} + +func (m *messageModel) toDomain() Message { + return Message{ + MessageBody: MessageBody{ + ExtID: m.ExtID, + Type: m.Type, + Sender: m.Sender, + Recipient: m.Recipient, + SimNumber: m.SimNumber, + Content: m.Content, + IsEncrypted: m.IsEncrypted, + CreatedAt: m.CreatedAt, + }, + + ID: m.ID, + DeviceID: m.DeviceID, + Attachments: lo.Map( + m.Attachments, + func(item attachmentModel, _ int) Attachment { return item.toDomain() }, + ), + } +} + +type attachmentModel struct { + ID uint64 `gorm:"primaryKey;type:BIGINT UNSIGNED;autoIncrement"` + MessageID uint64 `gorm:"not null;type:BIGINT UNSIGNED;uniqueIndex:unq_inbox_att_msg_part,priority:1"` + PartID int64 `gorm:"not null;type:BIGINT;uniqueIndex:unq_inbox_att_msg_part,priority:2"` + ContentType string `gorm:"not null;type:varchar(128)"` + Name string `gorm:"not null;type:varchar(512)"` + Size *int64 `gorm:"type:BIGINT"` + Data []byte `gorm:"not null;type:LONGBLOB"` +} + +func (*attachmentModel) TableName() string { + return "inbox_attachments" +} + +func (m *attachmentModel) toDomain() Attachment { + return Attachment{ + AttachmentInput: AttachmentInput{ + PartID: m.PartID, + ContentType: m.ContentType, + Name: m.Name, + Size: m.Size, + Data: m.Data, + }, + + ID: m.ID, + } +} + +func Migrate(db *gorm.DB) error { + if err := db.AutoMigrate(new(messageModel), new(attachmentModel)); err != nil { + return fmt.Errorf("inbox migration failed: %w", err) + } + return nil +} diff --git a/internal/sms-gateway/inbox/module.go b/internal/sms-gateway/inbox/module.go index 1feefbb8..c0d68322 100644 --- a/internal/sms-gateway/inbox/module.go +++ b/internal/sms-gateway/inbox/module.go @@ -1,16 +1,22 @@ package inbox import ( + "github.com/capcom6/go-infra-fx/db" "github.com/go-core-fx/logger" "go.uber.org/fx" ) +// Module returns the Fx module for the inbox messages feature. func Module() fx.Option { return fx.Module( "inbox", logger.WithNamedLogger("inbox"), - fx.Provide( - New, - ), + fx.Provide(NewRepository, fx.Private), + fx.Provide(NewService), ) } + +//nolint:gochecknoinits // framework-specific +func init() { + db.RegisterMigration(Migrate) +} diff --git a/internal/sms-gateway/inbox/repository.go b/internal/sms-gateway/inbox/repository.go new file mode 100644 index 00000000..ed6b2570 --- /dev/null +++ b/internal/sms-gateway/inbox/repository.go @@ -0,0 +1,149 @@ +package inbox + +import ( + "context" + "errors" + "fmt" + "time" + + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +// Repository provides data access for inbox messages. +type Repository struct { + db *gorm.DB +} + +// NewRepository creates a new Repository. +func NewRepository(db *gorm.DB) *Repository { + return &Repository{db: db} +} + +// InsertBatch stores multiple inbox messages in a single transaction. +// Uses OnConflict{DoNothing: true} for idempotency on the (ext_id, device_id) unique index. +func (r *Repository) InsertBatch(ctx context.Context, deviceID string, msgs []MessageInput) error { + if len(msgs) == 0 { + return nil + } + + models := make([]*messageModel, 0, len(msgs)) + for _, msg := range msgs { + models = append(models, newMessageModel(deviceID, msg)) + } + + err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + for _, msg := range models { + if err := insertMessageWithAttachments(tx, msg); err != nil { + return err + } + } + return nil + }) + if err != nil { + return fmt.Errorf("failed to insert inbox message batch: %w", err) + } + return nil +} + +func insertMessageWithAttachments(tx *gorm.DB, msg *messageModel) error { + result := tx.Omit("Device", "Attachments").Clauses( + clause.OnConflict{DoNothing: true}, + ).Create(msg) + if result.Error != nil { + return fmt.Errorf("failed to insert inbox message: %w", result.Error) + } + if result.RowsAffected == 0 { + return nil + } + if len(msg.Attachments) > 0 { + for i := range msg.Attachments { + msg.Attachments[i].MessageID = msg.ID + } + if err := tx.Create(&msg.Attachments).Error; err != nil { + return fmt.Errorf("failed to insert attachments: %w", err) + } + } + return nil +} + +// list returns inbox messages for a user with filtering and pagination. +func (r *Repository) list(userID string, filter ListFilter, opts ListOptions) ([]Message, int64, error) { + query := r.db.Model((*messageModel)(nil)). + Joins("JOIN devices ON inbox.device_id = devices.id"). + Where("devices.user_id = ?", userID) + + query = filter.apply(query) + + var total int64 + if err := query.Count(&total).Error; err != nil { + return nil, 0, fmt.Errorf("failed to count inbox messages: %w", err) + } + + query = opts.apply(query) + + query = query.Order("inbox.created_at DESC, inbox.id DESC") + + var messages []messageModel + if err := query.Preload("Attachments").Find(&messages).Error; err != nil { + return nil, 0, fmt.Errorf("failed to list inbox messages: %w", err) + } + + result := make([]Message, 0, len(messages)) + for i := range messages { + result = append(result, messages[i].toDomain()) + } + + return result, total, nil +} + +// findMessageByExtID finds a message by external ID, verifying user ownership via device join. +func (r *Repository) findMessageByExtID( + ctx context.Context, + extID string, + userID string, +) (*Message, error) { + var msg messageModel + err := r.db.WithContext(ctx). + Joins("JOIN devices ON inbox.device_id = devices.id"). + Where("inbox.ext_id = ? AND devices.user_id = ?", extID, userID). + Preload("Attachments"). + First(&msg).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("failed to find inbox message: %w", err) + } + domain := msg.toDomain() + return &domain, nil +} + +// findAttachment returns a single attachment by message ID and part ID. +func (r *Repository) findAttachment( + ctx context.Context, + messageID uint64, + partID int64, +) (*Attachment, error) { + var att attachmentModel + err := r.db.WithContext(ctx). + Where("message_id = ? AND part_id = ?", messageID, partID). + First(&att).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("failed to find attachment: %w", err) + } + domain := att.toDomain() + return &domain, nil +} + +// Cleanup deletes inbox messages older than the given time. +// Cascades to inbox_attachments via foreign key. +func (r *Repository) Cleanup(ctx context.Context, until time.Time) (int64, error) { + res := r.db.WithContext(ctx). + Where("created_at < ?", until). + Delete(new(messageModel)) + return res.RowsAffected, res.Error +} diff --git a/internal/sms-gateway/inbox/service.go b/internal/sms-gateway/inbox/service.go index 4eeb796a..537177c6 100644 --- a/internal/sms-gateway/inbox/service.go +++ b/internal/sms-gateway/inbox/service.go @@ -1,6 +1,7 @@ package inbox import ( + "context" "fmt" "time" @@ -9,16 +10,31 @@ import ( "go.uber.org/zap" ) +// Service provides business logic for inbox messages. type Service struct { + config Config + eventsSvc *events.Service + inbox *Repository + logger *zap.Logger } -func New(eventsSvc *events.Service, logger *zap.Logger) *Service { +// NewService creates a new Service. +func NewService( + config Config, + eventsSvc *events.Service, + inbox *Repository, + logger *zap.Logger, +) *Service { return &Service{ + config: config, + eventsSvc: eventsSvc, + inbox: inbox, + logger: logger, } } @@ -38,3 +54,44 @@ func (s *Service) Refresh( return nil } + +// InsertBatch stores a batch of encrypted inbox messages. +// Returns ErrNotEncrypted if any message has IsEncrypted=false. +// Returns ErrEmptyBatch if the slice is empty. +func (s *Service) InsertBatch(ctx context.Context, deviceID string, msgs []MessageInput) error { + if len(msgs) == 0 { + return ErrEmptyBatch + } + + for _, m := range msgs { + if !m.IsEncrypted { + return ErrNotEncrypted + } + } + + if err := s.inbox.InsertBatch(ctx, deviceID, msgs); err != nil { + return fmt.Errorf("failed to insert inbox messages: %w", err) + } + + return nil +} + +// List returns inbox messages for a user. +func (s *Service) List(userID string, filter ListFilter, opts ListOptions) ([]Message, int64, error) { + return s.inbox.list(userID, filter, opts) +} + +// GetAttachment returns a single attachment, verifying user ownership of the parent message. +func (s *Service) GetAttachment( + ctx context.Context, + userID string, + messageExtID string, + partID int64, +) (*Attachment, error) { + msg, err := s.inbox.findMessageByExtID(ctx, messageExtID, userID) + if err != nil { + return nil, err + } + + return s.inbox.findAttachment(ctx, msg.ID, partID) +} diff --git a/internal/sms-gateway/models/migrations/mysql/20260724000000_create_inbox_tables.sql b/internal/sms-gateway/models/migrations/mysql/20260724000000_create_inbox_tables.sql new file mode 100644 index 00000000..8b51f3d4 --- /dev/null +++ b/internal/sms-gateway/models/migrations/mysql/20260724000000_create_inbox_tables.sql @@ -0,0 +1,44 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE `inbox` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `ext_id` VARCHAR(36) NOT NULL, + `device_id` CHAR(21) NOT NULL, + `type` ENUM('SMS', 'DATA_SMS', 'MMS', 'MMS_DOWNLOADED') NOT NULL, + `sender` VARCHAR(512) NOT NULL, + `recipient` VARCHAR(512) NULL, + `sim_number` TINYINT(1) UNSIGNED NULL, + `content` TEXT NOT NULL, + `is_encrypted` TINYINT(1) UNSIGNED NOT NULL DEFAULT 1, + `created_at` DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), + `updated_at` DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3), + PRIMARY KEY (`id`), + UNIQUE INDEX `unq_inbox_ext_device` (`ext_id`, `device_id`), + INDEX `idx_inbox_device_created` (`device_id`, `created_at`), + INDEX `idx_inbox_created` (`created_at`), + CONSTRAINT `fk_inbox_device` FOREIGN KEY (`device_id`) REFERENCES `devices`(`id`) ON DELETE CASCADE +) ENGINE = InnoDB; +-- +goose StatementEnd +-- +goose StatementBegin +CREATE TABLE `inbox_attachments` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `message_id` BIGINT UNSIGNED NOT NULL, + `part_id` BIGINT NOT NULL, + `content_type` VARCHAR(128) NOT NULL, + `name` VARCHAR(512) NOT NULL, + `size` BIGINT NULL, + `data` LONGBLOB NOT NULL, + PRIMARY KEY (`id`), + UNIQUE INDEX `unq_inbox_att_msg_part` (`message_id`, `part_id`), + INDEX `idx_inbox_att_message` (`message_id`), + CONSTRAINT `fk_inbox_att_message` FOREIGN KEY (`message_id`) REFERENCES `inbox`(`id`) ON DELETE CASCADE +) ENGINE = InnoDB; +-- +goose StatementEnd +--- +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS `inbox_attachments`; +-- +goose StatementEnd +-- +goose StatementBegin +DROP TABLE IF EXISTS `inbox`; +-- +goose StatementEnd \ No newline at end of file diff --git a/internal/sms-gateway/openapi/docs.go b/internal/sms-gateway/openapi/docs.go index bee63823..f68eb1d3 100644 --- a/internal/sms-gateway/openapi/docs.go +++ b/internal/sms-gateway/openapi/docs.go @@ -362,7 +362,7 @@ const docTemplate = `{ "JWTAuth": [] } ], - "description": "Retrieves incoming messages with filtering and pagination.", + "description": "Retrieves inbox messages with filtering and pagination.", "produces": [ "application/json" ], @@ -370,7 +370,7 @@ const docTemplate = `{ "User", "Inbox" ], - "summary": "Get incoming messages", + "summary": "Get inbox messages", "parameters": [ { "enum": [ @@ -380,7 +380,7 @@ const docTemplate = `{ "MMS_DOWNLOADED" ], "type": "string", - "description": "Filter incoming messages by type", + "description": "Filter inbox messages by type", "name": "type", "in": "query" }, @@ -424,7 +424,7 @@ const docTemplate = `{ ], "responses": { "200": { - "description": "A list of incoming messages", + "description": "A list of inbox messages", "schema": { "type": "array", "items": { @@ -535,6 +535,81 @@ const docTemplate = `{ } } }, + "/3rdparty/v1/inbox/{id}/attachments/{partId}": { + "get": { + "security": [ + { + "ApiAuth": [] + }, + { + "JWTAuth": [] + } + ], + "description": "Downloads an attachment from an inbox message by message ID and part ID.", + "produces": [ + "application/octet-stream" + ], + "tags": [ + "User", + "Inbox" + ], + "summary": "Get attachment", + "parameters": [ + { + "type": "string", + "description": "Inbox message ID", + "name": "id", + "in": "path", + "required": true + }, + { + "type": "integer", + "description": "Attachment part ID", + "name": "partId", + "in": "path", + "required": true + } + ], + "responses": { + "200": { + "description": "Attachment file", + "schema": { + "type": "file" + } + }, + "400": { + "description": "Invalid request", + "schema": { + "$ref": "#/definitions/smsgateway.ErrorResponse" + } + }, + "401": { + "description": "Unauthorized", + "schema": { + "$ref": "#/definitions/smsgateway.ErrorResponse" + } + }, + "403": { + "description": "Forbidden", + "schema": { + "$ref": "#/definitions/smsgateway.ErrorResponse" + } + }, + "404": { + "description": "Not found", + "schema": { + "$ref": "#/definitions/smsgateway.ErrorResponse" + } + }, + "500": { + "description": "Internal server error", + "schema": { + "$ref": "#/definitions/smsgateway.ErrorResponse" + } + } + } + } + }, "/3rdparty/v1/logs": { "get": { "security": [ @@ -1699,6 +1774,31 @@ const docTemplate = `{ "HealthStatusFail" ] }, + "smsgateway.InboxAttachment": { + "type": "object", + "properties": { + "contentType": { + "description": "Attachment MIME type", + "type": "string", + "example": "image/jpeg" + }, + "name": { + "description": "Attachment file name", + "type": "string", + "example": "photo.jpg" + }, + "partId": { + "description": "Attachment part ID", + "type": "integer", + "example": 1 + }, + "size": { + "description": "Attachment file size", + "type": "integer", + "example": 102400 + } + } + }, "smsgateway.InboxRefreshRequest": { "type": "object", "required": [ @@ -1763,6 +1863,13 @@ const docTemplate = `{ "type" ], "properties": { + "attachments": { + "description": "MMS attachments", + "type": "array", + "items": { + "$ref": "#/definitions/smsgateway.InboxAttachment" + } + }, "contentPreview": { "description": "Message body preview or metadata", "type": "string", @@ -1775,17 +1882,22 @@ const docTemplate = `{ "example": "2020-01-01T00:00:00Z" }, "id": { - "description": "Incoming message ID", + "description": "Inbox message ID", "type": "string", "example": "PyDmBQZZXYmyxMwED8Fzy" }, + "isEncrypted": { + "description": "Whether the message is encrypted", + "type": "boolean", + "example": true + }, "recipient": { "description": "Recipient phone number on the device", "type": "string", "example": "+79990001234" }, "sender": { - "description": "Incoming sender phone number", + "description": "Inbox sender phone number", "type": "string", "example": "+79990001234" }, @@ -1839,6 +1951,7 @@ const docTemplate = `{ "devices:delete", "inbox:list", "inbox:refresh", + "inbox:read", "logs:read", "messages:cancel", "messages:send", @@ -1857,6 +1970,7 @@ const docTemplate = `{ "ScopeDevicesDelete", "ScopeInboxList", "ScopeInboxRefresh", + "ScopeInboxRead", "ScopeLogsRead", "ScopeMessagesCancel", "ScopeMessagesSend", diff --git a/internal/worker/config/config.go b/internal/worker/config/config.go index 7945615f..5c894c25 100644 --- a/internal/worker/config/config.go +++ b/internal/worker/config/config.go @@ -17,6 +17,7 @@ type Tasks struct { MessagesCleanup MessagesCleanup `yaml:"messages_cleanup"` DevicesCleanup DevicesCleanup `yaml:"devices_cleanup"` TokensCleanup TokensCleanup `yaml:"tokens_cleanup"` + InboxCleanup InboxCleanup `yaml:"inbox_cleanup"` } type MessagesHashing struct { Interval Duration `yaml:"interval" envconfig:"TASKS__MESSAGES_HASHING__INTERVAL"` @@ -37,6 +38,11 @@ type TokensCleanup struct { MaxAge Duration `yaml:"max_age" envconfig:"TASKS__TOKENS_CLEANUP__MAX_AGE"` } +type InboxCleanup struct { + Interval Duration `yaml:"interval" envconfig:"TASKS__INBOX_CLEANUP__INTERVAL"` + MaxAge Duration `yaml:"max_age" envconfig:"TASKS__INBOX_CLEANUP__MAX_AGE"` +} + func Default() Config { //nolint:exhaustruct,mnd,goconst // default values return Config{ @@ -56,6 +62,10 @@ func Default() Config { Interval: Duration(24 * time.Hour), MaxAge: Duration(1 * time.Hour), }, + InboxCleanup: InboxCleanup{ + Interval: Duration(24 * time.Hour), + MaxAge: Duration(30 * 24 * time.Hour), + }, }, Database: config.Database{ Host: "localhost", diff --git a/internal/worker/config/module.go b/internal/worker/config/module.go index ccf9b4d8..8d32b70e 100644 --- a/internal/worker/config/module.go +++ b/internal/worker/config/module.go @@ -6,6 +6,7 @@ import ( "github.com/android-sms-gateway/server/internal/worker/server" "github.com/android-sms-gateway/server/internal/worker/tasks/devices" + "github.com/android-sms-gateway/server/internal/worker/tasks/inbox" "github.com/android-sms-gateway/server/internal/worker/tasks/messages" "github.com/android-sms-gateway/server/internal/worker/tasks/tokens" "github.com/capcom6/go-infra-fx/config" @@ -69,6 +70,14 @@ func Module() fx.Option { }, } }), + fx.Provide(func(cfg Config) inbox.Config { + return inbox.Config{ + Cleanup: inbox.CleanupConfig{ + Interval: time.Duration(cfg.Tasks.InboxCleanup.Interval), + MaxAge: time.Duration(cfg.Tasks.InboxCleanup.MaxAge), + }, + } + }), fx.Provide(func(cfg Config) server.Config { return server.Config{ Address: cfg.HTTP.Listen, diff --git a/internal/worker/tasks/inbox/cleanup.go b/internal/worker/tasks/inbox/cleanup.go new file mode 100644 index 00000000..dec19b7f --- /dev/null +++ b/internal/worker/tasks/inbox/cleanup.go @@ -0,0 +1,56 @@ +package inbox + +import ( + "context" + "fmt" + "time" + + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" + "github.com/android-sms-gateway/server/internal/worker/executor" + "go.uber.org/zap" +) + +type cleanupTask struct { + config CleanupConfig + inbox *inbox.Repository + + logger *zap.Logger +} + +func NewCleanupTask( + config CleanupConfig, + inbox *inbox.Repository, + logger *zap.Logger, +) executor.PeriodicTask { + return &cleanupTask{ + config: config, + inbox: inbox, + logger: logger, + } +} + +// Interval implements executor.PeriodicTask. +func (c *cleanupTask) Interval() time.Duration { + return c.config.Interval +} + +// Name implements executor.PeriodicTask. +func (c *cleanupTask) Name() string { + return "inbox:cleanup" +} + +// Run implements executor.PeriodicTask. +func (c *cleanupTask) Run(ctx context.Context) error { + rows, err := c.inbox.Cleanup(ctx, time.Now().Add(-c.config.MaxAge)) + if err != nil { + return fmt.Errorf("failed to cleanup inbox messages: %w", err) + } + + if rows > 0 { + c.logger.Info("cleaned up inbox messages", zap.Int64("rows", rows)) + } + + return nil +} + +var _ executor.PeriodicTask = (*cleanupTask)(nil) diff --git a/internal/worker/tasks/inbox/config.go b/internal/worker/tasks/inbox/config.go new file mode 100644 index 00000000..3c6fa8ab --- /dev/null +++ b/internal/worker/tasks/inbox/config.go @@ -0,0 +1,12 @@ +package inbox + +import "time" + +type Config struct { + Cleanup CleanupConfig +} + +type CleanupConfig struct { + Interval time.Duration + MaxAge time.Duration +} diff --git a/internal/worker/tasks/inbox/module.go b/internal/worker/tasks/inbox/module.go new file mode 100644 index 00000000..07453a79 --- /dev/null +++ b/internal/worker/tasks/inbox/module.go @@ -0,0 +1,22 @@ +package inbox + +import ( + "github.com/android-sms-gateway/server/internal/sms-gateway/inbox" + "github.com/android-sms-gateway/server/internal/worker/executor" + "github.com/go-core-fx/logger" + "go.uber.org/fx" +) + +func Module() fx.Option { + return fx.Module( + "inbox", + logger.WithNamedLogger("inbox"), + fx.Provide(func(c Config) CleanupConfig { + return c.Cleanup + }, fx.Private), + fx.Provide(inbox.NewRepository, fx.Private), + fx.Provide( + executor.AsWorkerTask(NewCleanupTask), + ), + ) +} diff --git a/internal/worker/tasks/module.go b/internal/worker/tasks/module.go index 97dcad13..ef7b854f 100644 --- a/internal/worker/tasks/module.go +++ b/internal/worker/tasks/module.go @@ -2,6 +2,7 @@ package tasks import ( "github.com/android-sms-gateway/server/internal/worker/tasks/devices" + "github.com/android-sms-gateway/server/internal/worker/tasks/inbox" "github.com/android-sms-gateway/server/internal/worker/tasks/messages" "github.com/android-sms-gateway/server/internal/worker/tasks/tokens" "github.com/go-core-fx/logger" @@ -15,5 +16,6 @@ func Module() fx.Option { messages.Module(), devices.Module(), tokens.Module(), + inbox.Module(), ) } From 796c19b619843a370ba40af6101d349c4ffe999a Mon Sep 17 00:00:00 2001 From: Aleksandr Soloshenko Date: Thu, 20 Aug 2026 11:26:11 +0700 Subject: [PATCH 2/2] [lint] ignore FCM Fid migration --- internal/sms-gateway/modules/push/fcm/client.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/sms-gateway/modules/push/fcm/client.go b/internal/sms-gateway/modules/push/fcm/client.go index d99a893a..a862421d 100644 --- a/internal/sms-gateway/modules/push/fcm/client.go +++ b/internal/sms-gateway/modules/push/fcm/client.go @@ -69,7 +69,7 @@ func (c *Client) Send(ctx context.Context, messages []client.Message) ([]error, Android: &messaging.AndroidConfig{ Priority: "high", }, - Token: message.Token, + Token: message.Token, //nolint:staticcheck //migration to Fid is planned }) if err != nil { errs[i] = fmt.Errorf("failed to send message: %w", err)