From 0cd71c55064f606ec2c75395ae7096af4b71a487 Mon Sep 17 00:00:00 2001 From: phoenix Date: Thu, 13 Nov 2025 01:59:39 +0000 Subject: [PATCH] tsk-34: Fetch next ready scheduled message (#38) Closes #34 Reviewed-on: https://git.kundeng.us/phoenix/textsender-api/pulls/38 Co-authored-by: phoenix Co-committed-by: phoenix --- cmd/api/main.go | 1 + cmd/api/main_test.go | 1 + internal/handler/endpoint/endpoint.go | 1 + internal/handler/scheduled_message.go | 39 +++++++++++--- .../handler/scheduled_message_event_test.go | 1 - internal/handler/scheduled_message_test.go | 51 +++++++++++++++++++ .../store/mock/scheduled_message_store.go | 24 +++++++++ internal/store/scheduled_message_store.go | 28 ++++++++++ 8 files changed, 139 insertions(+), 7 deletions(-) diff --git a/cmd/api/main.go b/cmd/api/main.go index 7d0aa7e..611dfe0 100644 --- a/cmd/api/main.go +++ b/cmd/api/main.go @@ -89,6 +89,7 @@ func main() { router.Method("DELETE", endpoint.DeleteScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.DeleteScheduledMessageEvent))) router.Method("PATCH", endpoint.UpdateScheduledMessageStatusEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageStatusHandler.UpdateStatus))) router.Method("GET", endpoint.GetScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.GetScheduledMessage))) + router.Method("GET", endpoint.FetchNextScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.FetchNextMessage))) // Start server server := &http.Server{ diff --git a/cmd/api/main_test.go b/cmd/api/main_test.go index 57c9d84..3168e00 100644 --- a/cmd/api/main_test.go +++ b/cmd/api/main_test.go @@ -75,6 +75,7 @@ func TestMain(m *testing.M) { testRouter.Method("DELETE", endpoint.DeleteScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.DeleteScheduledMessageEvent))) testRouter.Method("PATCH", endpoint.UpdateScheduledMessageStatusEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageStatusHandler.UpdateStatus))) testRouter.Method("GET", endpoint.GetScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.GetScheduledMessage))) + testRouter.Method("GET", endpoint.FetchNextScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.FetchNextMessage))) code := m.Run() os.Exit(code) diff --git a/internal/handler/endpoint/endpoint.go b/internal/handler/endpoint/endpoint.go index 539381c..6160438 100644 --- a/internal/handler/endpoint/endpoint.go +++ b/internal/handler/endpoint/endpoint.go @@ -11,3 +11,4 @@ const AddEventToScheduledMessageEndpoint = "/api/v1/schedule/message/event" const GetScheduledMessageEventEndpoint = "/api/v1/schedule/message/event" const DeleteScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}" const UpdateScheduledMessageStatusEndpoint = "/api/v1/schedule/message/status/update" +const FetchNextScheduledMessageEndpoint = "/api/v1/schedule/message/fetch" diff --git a/internal/handler/scheduled_message.go b/internal/handler/scheduled_message.go index 8ccb601..73c1e98 100644 --- a/internal/handler/scheduled_message.go +++ b/internal/handler/scheduled_message.go @@ -29,6 +29,11 @@ type GetScheduledMessageResponse struct { Data []scheduling.ScheduledMessage `json:"data"` } +type FetchNextMessageResponse struct { + Message string `json:"message"` + Data []scheduling.ScheduledMessage `json:"data"` +} + type ScheduledMessageHandler struct { ScheduledMessageStore store.ScheduledMessageStore } @@ -37,7 +42,7 @@ func NewScheduledMessageHandler(str store.ScheduledMessageStore) *ScheduledMessa return &ScheduledMessageHandler{ScheduledMessageStore: str} } -func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *http.Request) { +func (s *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) return @@ -63,7 +68,7 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r * if validStatus, err := IsStatusValid(&scheduledMessage); err == nil { if validStatus { if valid, err := isScheduledTimeValid(&scheduledMessage); err == nil && valid { - if err = c.ScheduledMessageStore.CreateScheduledMessage(ctx, &scheduledMessage); err == nil { + if err = s.ScheduledMessageStore.CreateScheduledMessage(ctx, &scheduledMessage); err == nil { statusCode = http.StatusCreated resp.Data = append(resp.Data, scheduledMessage) resp.Message = "Successful" @@ -88,7 +93,7 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r * RespondWithJSON(w, statusCode, &resp) } -func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *http.Request) { +func (s *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *http.Request) { var id, userId uuid.UUID if idParam, err := ParseQueryParams(r, "id"); err == nil { @@ -115,7 +120,7 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r * http.Error(w, "Invalid query parameters", http.StatusBadRequest) return } else if userId != uuid.Nil && id != uuid.Nil { - if schMsgs, err := c.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil { + if schMsgs, err := s.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil { statusCode = http.StatusInternalServerError resp.Message = err.Error() } else { @@ -128,7 +133,7 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r * resp.Message = "Successful" } } else if id != uuid.Nil { - if schMsg, err := c.ScheduledMessageStore.Get(ctx, id); err != nil { + if schMsg, err := s.ScheduledMessageStore.Get(ctx, id); err != nil { statusCode = http.StatusInternalServerError resp.Message = err.Error() } else { @@ -137,7 +142,7 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r * resp.Data = append(resp.Data, *schMsg) } } else { - if schMsgs, err := c.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil { + if schMsgs, err := s.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil { statusCode = http.StatusInternalServerError resp.Message = err.Error() } else { @@ -152,6 +157,28 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r * RespondWithJSON(w, statusCode, &resp) } +func (s *ScheduledMessageHandler) FetchNextMessage(w http.ResponseWriter, r *http.Request) { + ctx := r.Context() + var resp FetchNextMessageResponse + var statusCode int + + if schMsg, err := s.ScheduledMessageStore.FetchNextScheduledMessage(ctx); err != nil { + statusCode = http.StatusInternalServerError + resp.Message = err.Error() + } else { + if schMsg != nil { + statusCode = http.StatusOK + resp.Message = "Successful" + resp.Data = append(resp.Data, *schMsg) + } else { + statusCode = http.StatusNotFound + resp.Message = "No item" + } + } + + RespondWithJSON(w, statusCode, &resp) +} + func isScheduledTimeValid(schMsg *scheduling.ScheduledMessage) (bool, error) { now := time.Now() timeCutOff := now.Add(-5 * time.Minute) diff --git a/internal/handler/scheduled_message_event_test.go b/internal/handler/scheduled_message_event_test.go index 579b419..bc17adb 100644 --- a/internal/handler/scheduled_message_event_test.go +++ b/internal/handler/scheduled_message_event_test.go @@ -115,7 +115,6 @@ func TestGetScheduledMessageEventWithMock(t *testing.T) { } url := fmt.Sprintf("%s?id=%s", endpoint.GetScheduledMessageEventEndpoint, event.Id.String()) - fmt.Println("Url:", url) req, _ := http.NewRequest("GET", url, nil) rr := httptest.NewRecorder() diff --git a/internal/handler/scheduled_message_test.go b/internal/handler/scheduled_message_test.go index 07a1849..ebfa145 100644 --- a/internal/handler/scheduled_message_test.go +++ b/internal/handler/scheduled_message_test.go @@ -78,6 +78,57 @@ func TestGetScheduledMessageWithMock(t *testing.T) { assert.NoError(t, err, "Error parsing response %v", err) } +func TestFetchScheduledMessageWithMock(t *testing.T) { + now := time.Now() + mockStore := mock.NewMockScheduledMessageStore() + contactStore := mock.NewMockContactStore() + messageStore := mock.NewMockMessageStore() + schMsgEventStore := mock.NewMockScheduledMessageEventStore() + + recipientId := uuid.New() + messageId := uuid.New() + scheduledMessageId := uuid.New() + testUserId := uuid.New() + + con := testContact(recipientId, testUserId) + msg := testMessage(messageId, testUserId) + schMsg := testScheduledMessage(scheduledMessageId, testUserId, now) + event := testScheduledMessageEvent(msg.Id, con.Id, schMsg.Id) + + ctx := t.Context() + + if err := contactStore.CreateContact(ctx, &con); err != nil { + assert.NoError(t, err, "Error creating contact: %v", err) + } else if err = messageStore.CreateMessage(ctx, &msg); err != nil { + assert.NoError(t, err, "Error creating message: %v", err) + } else if err = mockStore.CreateScheduledMessage(ctx, &schMsg); err != nil { + assert.NoError(t, err, "Error creating scheduled message: %v", err) + } else if err = schMsgEventStore.CreateScheduledMessageEvent(ctx, &event); err != nil { + assert.NoError(t, err, "Error creating scheduled message event: %v", err) + } else if newStatus, err := mockStore.UpdateStatus(ctx, schMsg.Id, scheduling.Ready); err != nil { + assert.NoError(t, err, "Error updating status: %v", err) + } else { + assert.Equal(t, *newStatus, scheduling.Ready) + } + + req, _ := http.NewRequest("GET", endpoint.FetchNextScheduledMessageEndpoint, nil) + rr := httptest.NewRecorder() + + handler := NewScheduledMessageHandler(mockStore) + handler.FetchNextMessage(rr, req) + + assert.Equal(t, http.StatusOK, rr.Code) + + var response FetchNextMessageResponse + err := json.Unmarshal(rr.Body.Bytes(), &response) + assert.NoError(t, err, "Error fetching scheduled message: %v", err) + + assert.NotEmpty(t, response.Data, "No data") + + fetchedSchMsg := response.Data[0] + assert.Equal(t, scheduling.Processing, fetchedSchMsg.Status, "The statuses do not match") +} + func testCreateScheduledMessageRequest(userId uuid.UUID, now time.Time) CreateScheduledMessageRequest { scheduled := now.Add(5 * time.Minute) return CreateScheduledMessageRequest{Status: scheduling.Pending, UserId: userId, Scheduled: scheduled} diff --git a/internal/store/mock/scheduled_message_store.go b/internal/store/mock/scheduled_message_store.go index 71adbf4..8211924 100644 --- a/internal/store/mock/scheduled_message_store.go +++ b/internal/store/mock/scheduled_message_store.go @@ -47,6 +47,30 @@ func (m *MockScheduledMessageStore) Get(ctx context.Context, id uuid.UUID) (*sch } } +func (m *MockScheduledMessageStore) FetchNextScheduledMessage(ctx context.Context) (*scheduling.ScheduledMessage, error) { + m.mu.Lock() + defer m.mu.Unlock() + + if m.Error != nil { + return nil, m.Error + } + + var schMsg *scheduling.ScheduledMessage + for _, msg := range m.ScheduledMessages { + if msg.Status == scheduling.Ready { + msg.Status = scheduling.Processing + schMsg = msg + break + } + } + + if schMsg == nil { + return nil, fmt.Errorf("No scheduled message is ready") + } else { + return schMsg, nil + } +} + func (m *MockScheduledMessageStore) GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error) { m.mu.Lock() defer m.mu.Unlock() diff --git a/internal/store/scheduled_message_store.go b/internal/store/scheduled_message_store.go index 02c8d33..a278135 100644 --- a/internal/store/scheduled_message_store.go +++ b/internal/store/scheduled_message_store.go @@ -13,6 +13,7 @@ import ( type ScheduledMessageStore interface { Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessage, error) + FetchNextScheduledMessage(ctx context.Context) (*scheduling.ScheduledMessage, error) GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error) CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error UpdateStatus(ctx context.Context, id uuid.UUID, updatedStatus string) (*string, error) @@ -45,6 +46,33 @@ func (s *PGScheduledMessageStore) Get(ctx context.Context, id uuid.UUID) (*sched return &schMsg, nil } +func (s *PGScheduledMessageStore) FetchNextScheduledMessage(ctx context.Context) (*scheduling.ScheduledMessage, error) { + query := ` + UPDATE scheduled_messages + SET status = $1 + WHERE id = ( + SELECT id FROM scheduled_messages + WHERE status = $2 + ORDER BY id + FOR UPDATE SKIP LOCKED + LIMIT 1 + ) + RETURNING id, scheduled, created, status, user_id + ` + var schMsg scheduling.ScheduledMessage + err := s.db.QueryRow(ctx, query, scheduling.Processing, scheduling.Ready).Scan( + &schMsg.Id, &schMsg.Scheduled, &schMsg.Created, &schMsg.Status, &schMsg.UserId, + ) + + if err == pgx.ErrNoRows { + return nil, nil + } else if err != nil { + return nil, fmt.Errorf("Getting scheduled message: %w", err) + } else { + return &schMsg, nil + } +} + func (s *PGScheduledMessageStore) GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error) { query := ` SELECT id, scheduled, created, status, user_id FROM scheduled_messages WHERE user_id = $1