Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0cd71c5506 | ||
|
|
3218c69f1e |
@@ -89,6 +89,7 @@ func main() {
|
|||||||
router.Method("DELETE", endpoint.DeleteScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.DeleteScheduledMessageEvent)))
|
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("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.GetScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.GetScheduledMessage)))
|
||||||
|
router.Method("GET", endpoint.FetchNextScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.FetchNextMessage)))
|
||||||
|
|
||||||
// Start server
|
// Start server
|
||||||
server := &http.Server{
|
server := &http.Server{
|
||||||
|
|||||||
@@ -75,6 +75,7 @@ func TestMain(m *testing.M) {
|
|||||||
testRouter.Method("DELETE", endpoint.DeleteScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.DeleteScheduledMessageEvent)))
|
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("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.GetScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.GetScheduledMessage)))
|
||||||
|
testRouter.Method("GET", endpoint.FetchNextScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.FetchNextMessage)))
|
||||||
|
|
||||||
code := m.Run()
|
code := m.Run()
|
||||||
os.Exit(code)
|
os.Exit(code)
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ const ADD_CONTACT_ENDPOINT = "/api/v1/contact/new"
|
|||||||
const ScheduleMessageEndpoint = "/api/v1/schedule/message"
|
const ScheduleMessageEndpoint = "/api/v1/schedule/message"
|
||||||
const GetScheduledMessageEndpoint = "/api/v1/schedule/message"
|
const GetScheduledMessageEndpoint = "/api/v1/schedule/message"
|
||||||
const AddEventToScheduledMessageEndpoint = "/api/v1/schedule/message/event"
|
const AddEventToScheduledMessageEndpoint = "/api/v1/schedule/message/event"
|
||||||
const GetScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}"
|
const GetScheduledMessageEventEndpoint = "/api/v1/schedule/message/event"
|
||||||
const DeleteScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}"
|
const DeleteScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}"
|
||||||
const UpdateScheduledMessageStatusEndpoint = "/api/v1/schedule/message/status/update"
|
const UpdateScheduledMessageStatusEndpoint = "/api/v1/schedule/message/status/update"
|
||||||
|
const FetchNextScheduledMessageEndpoint = "/api/v1/schedule/message/fetch"
|
||||||
|
|||||||
@@ -29,6 +29,11 @@ type GetScheduledMessageResponse struct {
|
|||||||
Data []scheduling.ScheduledMessage `json:"data"`
|
Data []scheduling.ScheduledMessage `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type FetchNextMessageResponse struct {
|
||||||
|
Message string `json:"message"`
|
||||||
|
Data []scheduling.ScheduledMessage `json:"data"`
|
||||||
|
}
|
||||||
|
|
||||||
type ScheduledMessageHandler struct {
|
type ScheduledMessageHandler struct {
|
||||||
ScheduledMessageStore store.ScheduledMessageStore
|
ScheduledMessageStore store.ScheduledMessageStore
|
||||||
}
|
}
|
||||||
@@ -37,7 +42,7 @@ func NewScheduledMessageHandler(str store.ScheduledMessageStore) *ScheduledMessa
|
|||||||
return &ScheduledMessageHandler{ScheduledMessageStore: str}
|
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 {
|
if r.Method != http.MethodPost {
|
||||||
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
|
||||||
return
|
return
|
||||||
@@ -63,7 +68,7 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *
|
|||||||
if validStatus, err := IsStatusValid(&scheduledMessage); err == nil {
|
if validStatus, err := IsStatusValid(&scheduledMessage); err == nil {
|
||||||
if validStatus {
|
if validStatus {
|
||||||
if valid, err := isScheduledTimeValid(&scheduledMessage); err == nil && valid {
|
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
|
statusCode = http.StatusCreated
|
||||||
resp.Data = append(resp.Data, scheduledMessage)
|
resp.Data = append(resp.Data, scheduledMessage)
|
||||||
resp.Message = "Successful"
|
resp.Message = "Successful"
|
||||||
@@ -88,7 +93,7 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *
|
|||||||
RespondWithJSON(w, statusCode, &resp)
|
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
|
var id, userId uuid.UUID
|
||||||
|
|
||||||
if idParam, err := ParseQueryParams(r, "id"); err == nil {
|
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)
|
http.Error(w, "Invalid query parameters", http.StatusBadRequest)
|
||||||
return
|
return
|
||||||
} else if userId != uuid.Nil && id != uuid.Nil {
|
} 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
|
statusCode = http.StatusInternalServerError
|
||||||
resp.Message = err.Error()
|
resp.Message = err.Error()
|
||||||
} else {
|
} else {
|
||||||
@@ -128,7 +133,7 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *
|
|||||||
resp.Message = "Successful"
|
resp.Message = "Successful"
|
||||||
}
|
}
|
||||||
} else if id != uuid.Nil {
|
} 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
|
statusCode = http.StatusInternalServerError
|
||||||
resp.Message = err.Error()
|
resp.Message = err.Error()
|
||||||
} else {
|
} else {
|
||||||
@@ -137,7 +142,7 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *
|
|||||||
resp.Data = append(resp.Data, *schMsg)
|
resp.Data = append(resp.Data, *schMsg)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
if schMsgs, err := c.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil {
|
if schMsgs, err := s.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil {
|
||||||
statusCode = http.StatusInternalServerError
|
statusCode = http.StatusInternalServerError
|
||||||
resp.Message = err.Error()
|
resp.Message = err.Error()
|
||||||
} else {
|
} else {
|
||||||
@@ -152,6 +157,28 @@ func (c *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *
|
|||||||
RespondWithJSON(w, statusCode, &resp)
|
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) {
|
func isScheduledTimeValid(schMsg *scheduling.ScheduledMessage) (bool, error) {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
timeCutOff := now.Add(-5 * time.Minute)
|
timeCutOff := now.Add(-5 * time.Minute)
|
||||||
|
|||||||
@@ -109,26 +109,23 @@ func (s *ScheduledMessageEventHandler) GetScheduledMessageEvent(w http.ResponseW
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
id := chi.URLParam(r, "id")
|
var id, scheduledMessageId uuid.UUID
|
||||||
if len(id) == 0 {
|
if idParam, err := ParseQueryParams(r, "id"); err == nil {
|
||||||
pathParts := strings.Split(r.URL.Path, "/")
|
id, err = uuid.Parse(*idParam)
|
||||||
if len(pathParts) < 7 {
|
}
|
||||||
http.Error(w, "Id not provided", http.StatusBadRequest)
|
if scheduledMessageIdParam, err := ParseQueryParams(r, "scheduled_message_id"); err == nil {
|
||||||
return
|
scheduledMessageId, err = uuid.Parse(*scheduledMessageIdParam)
|
||||||
} else {
|
|
||||||
id = pathParts[6]
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var statusCode int
|
var statusCode int
|
||||||
var resp GetScheduledMessageEventResponse
|
var resp GetScheduledMessageEventResponse
|
||||||
|
ctx := r.Context()
|
||||||
|
|
||||||
if parsedId, err := uuid.Parse(id); err != nil {
|
if id == uuid.Nil && scheduledMessageId == uuid.Nil {
|
||||||
resp.Message = err.Error()
|
|
||||||
statusCode = http.StatusBadRequest
|
statusCode = http.StatusBadRequest
|
||||||
} else {
|
resp.Message = "Query parameters missing"
|
||||||
ctx := r.Context()
|
} else if id != uuid.Nil {
|
||||||
if event, err := s.ScheduledMessageEventStore.Get(ctx, parsedId); err != nil {
|
if event, err := s.ScheduledMessageEventStore.Get(ctx, id); err != nil {
|
||||||
resp.Message = err.Error()
|
resp.Message = err.Error()
|
||||||
statusCode = http.StatusInternalServerError
|
statusCode = http.StatusInternalServerError
|
||||||
} else {
|
} else {
|
||||||
@@ -141,6 +138,17 @@ func (s *ScheduledMessageEventHandler) GetScheduledMessageEvent(w http.ResponseW
|
|||||||
resp.Message = "Scheduled message event not found"
|
resp.Message = "Scheduled message event not found"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
} else {
|
||||||
|
if events, err := s.ScheduledMessageEventStore.GetWithScheduleMessageId(ctx, scheduledMessageId); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
statusCode = http.StatusOK
|
||||||
|
resp.Message = "Successful"
|
||||||
|
for _, event := range events {
|
||||||
|
resp.Data = append(resp.Data, *event)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
RespondWithJSON(w, statusCode, &resp)
|
RespondWithJSON(w, statusCode, &resp)
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package handler
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -113,8 +114,8 @@ func TestGetScheduledMessageEventWithMock(t *testing.T) {
|
|||||||
assert.NoError(t, err, "Error creating scheduled message event: %v", err)
|
assert.NoError(t, err, "Error creating scheduled message event: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
endpointValue := strings.Replace(endpoint.GetScheduledMessageEventEndpoint, "{id}", event.Id.String(), 1)
|
url := fmt.Sprintf("%s?id=%s", endpoint.GetScheduledMessageEventEndpoint, event.Id.String())
|
||||||
req, _ := http.NewRequest("GET", endpointValue, nil)
|
req, _ := http.NewRequest("GET", url, nil)
|
||||||
rr := httptest.NewRecorder()
|
rr := httptest.NewRecorder()
|
||||||
|
|
||||||
handler.GetScheduledMessageEvent(rr, req)
|
handler.GetScheduledMessageEvent(rr, req)
|
||||||
|
|||||||
@@ -78,6 +78,57 @@ func TestGetScheduledMessageWithMock(t *testing.T) {
|
|||||||
assert.NoError(t, err, "Error parsing response %v", err)
|
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 {
|
func testCreateScheduledMessageRequest(userId uuid.UUID, now time.Time) CreateScheduledMessageRequest {
|
||||||
scheduled := now.Add(5 * time.Minute)
|
scheduled := now.Add(5 * time.Minute)
|
||||||
return CreateScheduledMessageRequest{Status: scheduling.Pending, UserId: userId, Scheduled: scheduled}
|
return CreateScheduledMessageRequest{Status: scheduling.Pending, UserId: userId, Scheduled: scheduled}
|
||||||
|
|||||||
@@ -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) {
|
func (m *MockScheduledMessageStore) GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error) {
|
||||||
m.mu.Lock()
|
m.mu.Lock()
|
||||||
defer m.mu.Unlock()
|
defer m.mu.Unlock()
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import (
|
|||||||
|
|
||||||
type ScheduledMessageStore interface {
|
type ScheduledMessageStore interface {
|
||||||
Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessage, error)
|
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)
|
GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error)
|
||||||
CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error
|
CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error
|
||||||
UpdateStatus(ctx context.Context, id uuid.UUID, updatedStatus string) (*string, 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
|
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) {
|
func (s *PGScheduledMessageStore) GetWithUserId(ctx context.Context, userId uuid.UUID) ([]*scheduling.ScheduledMessage, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, scheduled, created, status, user_id FROM scheduled_messages WHERE user_id = $1
|
SELECT id, scheduled, created, status, user_id FROM scheduled_messages WHERE user_id = $1
|
||||||
|
|||||||
Reference in New Issue
Block a user