tsk-19: Update status of scheduled message endpoint (#32)
Closes #19 Reviewed-on: phoenix/textsender-api#32 Co-authored-by: phoenix <kundeng00@pm.me> Co-committed-by: phoenix <kundeng00@pm.me>
This commit is contained in:
@@ -70,6 +70,7 @@ func main() {
|
|||||||
messageHandler := handler.NewMessageHandler(messageStore)
|
messageHandler := handler.NewMessageHandler(messageStore)
|
||||||
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
|
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
|
||||||
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore, schStore)
|
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore, schStore)
|
||||||
|
scheduledMessageStatusHandler := handler.NewScheduledMessageStatusHandler(schMsgEventStore, schStore)
|
||||||
|
|
||||||
router := chi.NewRouter()
|
router := chi.NewRouter()
|
||||||
|
|
||||||
@@ -86,6 +87,7 @@ func main() {
|
|||||||
router.Handle(endpoint.AddEventToScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.AddScheduledMessageEvent)))
|
router.Handle(endpoint.AddEventToScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.AddScheduledMessageEvent)))
|
||||||
router.Method("GET", endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent)))
|
router.Method("GET", endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent)))
|
||||||
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)))
|
||||||
|
|
||||||
// Start server
|
// Start server
|
||||||
server := &http.Server{
|
server := &http.Server{
|
||||||
|
|||||||
@@ -63,6 +63,7 @@ func TestMain(m *testing.M) {
|
|||||||
messageHandler := handler.NewMessageHandler(messageStore)
|
messageHandler := handler.NewMessageHandler(messageStore)
|
||||||
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
|
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
|
||||||
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore, schStore)
|
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore, schStore)
|
||||||
|
scheduledMessageStatusHandler := handler.NewScheduledMessageStatusHandler(schMsgEventStore, schStore)
|
||||||
|
|
||||||
testRouter = chi.NewRouter()
|
testRouter = chi.NewRouter()
|
||||||
testRouter.Handle(endpoint.ADD_CONTACT_ENDPOINT, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(contactHandler.AddContact)))
|
testRouter.Handle(endpoint.ADD_CONTACT_ENDPOINT, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(contactHandler.AddContact)))
|
||||||
@@ -72,6 +73,7 @@ func TestMain(m *testing.M) {
|
|||||||
testRouter.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage)))
|
testRouter.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage)))
|
||||||
testRouter.Method("GET", endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent)))
|
testRouter.Method("GET", endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent)))
|
||||||
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)))
|
||||||
|
|
||||||
code := m.Run()
|
code := m.Run()
|
||||||
os.Exit(code)
|
os.Exit(code)
|
||||||
|
|||||||
@@ -9,3 +9,4 @@ const ScheduleMessageEndpoint = "/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/{id}"
|
||||||
const DeleteScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}"
|
const DeleteScheduledMessageEventEndpoint = "/api/v1/schedule/message/event/{id}"
|
||||||
|
const UpdateScheduledMessageStatusEndpoint = "/api/v1/schedule/message/status/update"
|
||||||
|
|||||||
@@ -0,0 +1,154 @@
|
|||||||
|
package handler
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
|
||||||
|
"git.kundeng.us/phoenix/textsender-api/internal/store"
|
||||||
|
"git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling"
|
||||||
|
)
|
||||||
|
|
||||||
|
type RequestScheduledMessageStatus struct {
|
||||||
|
ScheduledMessageId uuid.UUID `json:"scheduled_message_id"`
|
||||||
|
Status string `json:"status"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ScheduledMessageChange struct {
|
||||||
|
OldStatus string `json:"old_status"`
|
||||||
|
ScheduledMessage scheduling.ScheduledMessage `json:"scheduled_message"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ScheduledMessageStatusResponse struct {
|
||||||
|
Message string `json:"message"`
|
||||||
|
Data []ScheduledMessageChange `json:"data"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ScheduledMessageStatusHandler struct {
|
||||||
|
ScheduledMessageEventStore store.ScheduledMessageEventStore
|
||||||
|
ScheduledMessageStore store.ScheduledMessageStore
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewScheduledMessageStatusHandler(str store.ScheduledMessageEventStore, schStore store.ScheduledMessageStore) *ScheduledMessageStatusHandler {
|
||||||
|
return &ScheduledMessageStatusHandler{ScheduledMessageEventStore: str, ScheduledMessageStore: schStore}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *ScheduledMessageStatusHandler) UpdateStatus(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Method != http.MethodPatch {
|
||||||
|
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
var req RequestScheduledMessageStatus
|
||||||
|
if err := ExtractFromRequest(r, &req); err != nil {
|
||||||
|
http.Error(w, "Invalid JSON: "+err.Error(), http.StatusBadRequest)
|
||||||
|
}
|
||||||
|
defer r.Body.Close()
|
||||||
|
|
||||||
|
var statusCode int
|
||||||
|
var resp ScheduledMessageStatusResponse
|
||||||
|
|
||||||
|
if req.ScheduledMessageId == uuid.Nil {
|
||||||
|
statusCode = http.StatusBadRequest
|
||||||
|
resp.Message = "Scheduled messaged Id is nil"
|
||||||
|
} else if len(req.Status) == 0 {
|
||||||
|
statusCode = http.StatusBadRequest
|
||||||
|
resp.Message = "No status provided"
|
||||||
|
} else {
|
||||||
|
ctx := r.Context()
|
||||||
|
|
||||||
|
if schMsg, err := s.ScheduledMessageStore.Get(ctx, req.ScheduledMessageId); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
if schMsg.Status == scheduling.Done {
|
||||||
|
statusCode = http.StatusNotFound
|
||||||
|
resp.Message = "Scheduled message has already completed"
|
||||||
|
} else {
|
||||||
|
switch req.Status {
|
||||||
|
case scheduling.Pending:
|
||||||
|
{
|
||||||
|
if schMsg.Status == scheduling.Processing {
|
||||||
|
statusCode = http.StatusBadRequest
|
||||||
|
resp.Message = "Message is currently processing"
|
||||||
|
} else {
|
||||||
|
if returnedStatus, err := s.ScheduledMessageStore.UpdateStatus(ctx, schMsg.Id, req.Status); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
chg := initScheduledMessageChange(schMsg, schMsg.Status, *returnedStatus)
|
||||||
|
statusCode = http.StatusOK
|
||||||
|
resp.Message = "Successful"
|
||||||
|
resp.Data = append(resp.Data, chg)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case scheduling.Ready:
|
||||||
|
{
|
||||||
|
if events, err := s.ScheduledMessageEventStore.GetWithScheduleMessageId(ctx, schMsg.Id); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
oldStatus := schMsg.Status
|
||||||
|
if schMsg.Status != scheduling.Processing && schMsg.Status != scheduling.Ready {
|
||||||
|
// Update status
|
||||||
|
if len(events) > 0 {
|
||||||
|
if returnedStatus, err := s.ScheduledMessageStore.UpdateStatus(ctx, schMsg.Id, req.Status); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
chg := initScheduledMessageChange(schMsg, oldStatus, *returnedStatus)
|
||||||
|
statusCode = http.StatusOK
|
||||||
|
resp.Message = "Successful"
|
||||||
|
resp.Data = append(resp.Data, chg)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = "Not enough scheduled message events"
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
if schMsg.Status == scheduling.Processing {
|
||||||
|
statusCode = http.StatusBadRequest
|
||||||
|
resp.Message = "Scheduled messages are processing"
|
||||||
|
} else {
|
||||||
|
statusCode = http.StatusNotModified
|
||||||
|
resp.Message = "Status is already set"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case scheduling.Processing, scheduling.Done:
|
||||||
|
{
|
||||||
|
if events, err := s.ScheduledMessageEventStore.GetWithScheduleMessageId(ctx, schMsg.Id); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
if len(events) > 0 {
|
||||||
|
if returnedStatus, err := s.ScheduledMessageStore.UpdateStatus(ctx, schMsg.Id, req.Status); err != nil {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = err.Error()
|
||||||
|
} else {
|
||||||
|
chg := initScheduledMessageChange(schMsg, schMsg.Status, *returnedStatus)
|
||||||
|
statusCode = http.StatusOK
|
||||||
|
resp.Message = "Successful"
|
||||||
|
resp.Data = append(resp.Data, chg)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
statusCode = http.StatusInternalServerError
|
||||||
|
resp.Message = "Not enough scheduled message events"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
RespondWithJSON(w, statusCode, &resp)
|
||||||
|
}
|
||||||
|
|
||||||
|
func initScheduledMessageChange(schMsg *scheduling.ScheduledMessage, oldStatus, newStatus string) ScheduledMessageChange {
|
||||||
|
schMsg.Status = newStatus
|
||||||
|
return ScheduledMessageChange{OldStatus: oldStatus, ScheduledMessage: *schMsg}
|
||||||
|
}
|
||||||
@@ -0,0 +1,73 @@
|
|||||||
|
package handler
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling"
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
|
||||||
|
"git.kundeng.us/phoenix/textsender-api/internal/handler/endpoint"
|
||||||
|
"git.kundeng.us/phoenix/textsender-api/internal/store/mock"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestUpdateScheduledMessageStatusWithMock(t *testing.T) {
|
||||||
|
now := time.Now()
|
||||||
|
|
||||||
|
schMsgEventStore := mock.NewMockScheduledMessageEventStore()
|
||||||
|
contactStore := mock.NewMockContactStore()
|
||||||
|
messageStore := mock.NewMockMessageStore()
|
||||||
|
schMsgStore := mock.NewMockScheduledMessageStore()
|
||||||
|
|
||||||
|
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 = schMsgStore.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)
|
||||||
|
}
|
||||||
|
|
||||||
|
testReq := RequestScheduledMessageStatus{}
|
||||||
|
testReq.Status = scheduling.Ready
|
||||||
|
testReq.ScheduledMessageId = schMsg.Id
|
||||||
|
jsonValue, _ := json.Marshal(testReq)
|
||||||
|
jsonBody := strings.NewReader(string(jsonValue))
|
||||||
|
|
||||||
|
req, _ := http.NewRequest("PATCH", endpoint.UpdateScheduledMessageStatusEndpoint, jsonBody)
|
||||||
|
rr := httptest.NewRecorder()
|
||||||
|
|
||||||
|
handler := NewScheduledMessageStatusHandler(schMsgEventStore, schMsgStore)
|
||||||
|
handler.UpdateStatus(rr, req)
|
||||||
|
|
||||||
|
assert.Equal(t, http.StatusOK, rr.Code)
|
||||||
|
|
||||||
|
var response ScheduledMessageStatusResponse
|
||||||
|
err := json.Unmarshal(rr.Body.Bytes(), &response)
|
||||||
|
assert.NoError(t, err, "Error creating event %v", err)
|
||||||
|
|
||||||
|
assert.NotEmpty(t, response.Data, "Data should not be empty")
|
||||||
|
|
||||||
|
changes := response.Data[0]
|
||||||
|
assert.NotEqual(t, changes.OldStatus, changes.ScheduledMessage.Status, "The status should not match")
|
||||||
|
assert.Equal(t, scheduling.Pending, changes.OldStatus, "The Old status does not match Old %s New %s", changes.OldStatus, scheduling.Pending)
|
||||||
|
assert.Equal(t, scheduling.Ready, changes.ScheduledMessage.Status, "Status has not been updated")
|
||||||
|
}
|
||||||
@@ -55,7 +55,7 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *
|
|||||||
} else {
|
} else {
|
||||||
ctx := r.Context()
|
ctx := r.Context()
|
||||||
|
|
||||||
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 = c.ScheduledMessageStore.CreateScheduledMessage(ctx, &scheduledMessage); err == nil {
|
||||||
@@ -99,7 +99,7 @@ func isScheduledTimeValid(schMsg *scheduling.ScheduledMessage) (bool, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func isStatusValid(schMsg *scheduling.ScheduledMessage) (bool, error) {
|
func IsStatusValid(schMsg *scheduling.ScheduledMessage) (bool, error) {
|
||||||
schMsg.Status = strings.ToUpper(schMsg.Status)
|
schMsg.Status = strings.ToUpper(schMsg.Status)
|
||||||
if schMsg.Status == scheduling.Pending {
|
if schMsg.Status == scheduling.Pending {
|
||||||
return true, nil
|
return true, nil
|
||||||
|
|||||||
@@ -50,6 +50,33 @@ func (m *MockScheduledMessageEventStore) Get(ctx context.Context, id uuid.UUID)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *MockScheduledMessageEventStore) GetWithScheduleMessageId(ctx context.Context, schMsgId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
if m.Error != nil {
|
||||||
|
return nil, m.Error
|
||||||
|
}
|
||||||
|
|
||||||
|
if schMsgId == uuid.Nil {
|
||||||
|
return nil, fmt.Errorf("Scheduled message Id is nil")
|
||||||
|
}
|
||||||
|
|
||||||
|
var events []*scheduling.ScheduledMessageEvent
|
||||||
|
|
||||||
|
for _, event := range m.ScheduledMessageEvents {
|
||||||
|
if event.ScheduledMessageId == schMsgId {
|
||||||
|
events = append(events, event)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(events) > 0 {
|
||||||
|
return events, nil
|
||||||
|
} else {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (m *MockScheduledMessageEventStore) Delete(ctx context.Context, id uuid.UUID) error {
|
func (m *MockScheduledMessageEventStore) Delete(ctx context.Context, id uuid.UUID) error {
|
||||||
m.mu.Lock()
|
m.mu.Lock()
|
||||||
defer m.mu.Unlock()
|
defer m.mu.Unlock()
|
||||||
|
|||||||
@@ -65,3 +65,28 @@ func (m *MockScheduledMessageStore) CreateScheduledMessage(ctx context.Context,
|
|||||||
m.ScheduledMessagesByKey[key] = schedMsg
|
m.ScheduledMessagesByKey[key] = schedMsg
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *MockScheduledMessageStore) UpdateStatus(ctx context.Context, id uuid.UUID, updatedStatus string) (*string, error) {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
if m.Error != nil {
|
||||||
|
return nil, m.Error
|
||||||
|
}
|
||||||
|
|
||||||
|
if id == uuid.Nil {
|
||||||
|
return nil, fmt.Errorf("Id is nil")
|
||||||
|
} else {
|
||||||
|
if schMsg := m.ScheduledMessages[id]; schMsg != nil {
|
||||||
|
key := ScheduledMessageKey{UserId: schMsg.UserId}
|
||||||
|
|
||||||
|
schMsg.Status = updatedStatus
|
||||||
|
m.ScheduledMessages[schMsg.Id] = schMsg
|
||||||
|
m.ScheduledMessagesByKey[key] = schMsg
|
||||||
|
|
||||||
|
return &updatedStatus, nil
|
||||||
|
} else {
|
||||||
|
return nil, fmt.Errorf("Scheduled message does not exist")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import (
|
|||||||
|
|
||||||
type ScheduledMessageEventStore interface {
|
type ScheduledMessageEventStore interface {
|
||||||
Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessageEvent, error)
|
Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessageEvent, error)
|
||||||
|
GetWithScheduleMessageId(ctx context.Context, schMsgId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error)
|
||||||
Delete(ctx context.Context, id uuid.UUID) error
|
Delete(ctx context.Context, id uuid.UUID) error
|
||||||
CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error
|
CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error
|
||||||
Exists(ctx context.Context, event *scheduling.ScheduledMessageEvent) (bool, error)
|
Exists(ctx context.Context, event *scheduling.ScheduledMessageEvent) (bool, error)
|
||||||
@@ -43,6 +44,34 @@ func (s *PGScheduledMessageEventStore) Get(ctx context.Context, id uuid.UUID) (*
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *PGScheduledMessageEventStore) GetWithScheduleMessageId(ctx context.Context, schMsgId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
|
||||||
|
var events []*scheduling.ScheduledMessageEvent
|
||||||
|
query := `SELECT id, recipient_id, message_id, scheduled_message_id, created FROM scheduled_message_events WHERE scheduled_message_id = $1
|
||||||
|
`
|
||||||
|
|
||||||
|
rows, err := s.db.Query(ctx, query, schMsgId)
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("Error retrieving rows from scheduled_message_Events: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for rows.Next() {
|
||||||
|
var event scheduling.ScheduledMessageEvent
|
||||||
|
|
||||||
|
if err := rows.Scan(&event.Id, &event.RecipientId, &event.MessageId, &event.ScheduledMessageId, &event.Created); err != nil {
|
||||||
|
return nil, fmt.Errorf("Error fetching data from row: %w", err)
|
||||||
|
} else {
|
||||||
|
events = append(events, &event)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := rows.Err(); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return events, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (s *PGScheduledMessageEventStore) Delete(ctx context.Context, id uuid.UUID) error {
|
func (s *PGScheduledMessageEventStore) Delete(ctx context.Context, id uuid.UUID) error {
|
||||||
query := `
|
query := `
|
||||||
DELETE FROM scheduled_message_events WHERE id = $1
|
DELETE FROM scheduled_message_events WHERE id = $1
|
||||||
|
|||||||
@@ -14,6 +14,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)
|
||||||
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)
|
||||||
}
|
}
|
||||||
|
|
||||||
type PGScheduledMessageStore struct {
|
type PGScheduledMessageStore struct {
|
||||||
@@ -54,3 +55,18 @@ func (s *PGScheduledMessageStore) CreateScheduledMessage(ctx context.Context, sc
|
|||||||
&schedMsg.Id, &schedMsg.Created,
|
&schedMsg.Id, &schedMsg.Created,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *PGScheduledMessageStore) UpdateStatus(ctx context.Context, id uuid.UUID, updatedStatus string) (*string, error) {
|
||||||
|
query := `
|
||||||
|
UPDATE scheduled_messages set status = $1 WHERE id = $2
|
||||||
|
RETURNING status
|
||||||
|
`
|
||||||
|
|
||||||
|
var returnedStatus string
|
||||||
|
|
||||||
|
if err := s.db.QueryRow(ctx, query, updatedStatus, id).Scan(&returnedStatus); err != nil {
|
||||||
|
return nil, err
|
||||||
|
} else {
|
||||||
|
return &returnedStatus, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user