6 Commits
Author SHA1 Message Date
phoenixandphoenix 0cd71c5506 tsk-34: Fetch next ready scheduled message (#38)
Closes #34

Reviewed-on: phoenix/textsender-api#38
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-13 01:59:39 +00:00
phoenixandphoenix 3218c69f1e tsk-33: Tweak get schedule message event endpoint (#37)
Closes #33

Reviewed-on: phoenix/textsender-api#37
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-13 00:48:04 +00:00
phoenixandphoenix 7f588cc80c tsk-13: Retrieve scheduled message endpoint (#36)
Closes #13

Reviewed-on: phoenix/textsender-api#36
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-12 19:25:54 +00:00
phoenixandphoenix a37f1627fc 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>
2025-11-12 17:54:32 +00:00
phoenixandphoenix d191d6a9c1 tsk-27: Added endpoint to delete scheduled message event (#31)
Closes #27

Reviewed-on: phoenix/textsender-api#31
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-11 18:35:04 +00:00
phoenixandphoenix 2716bb50d3 tsk-29: Place restrictions on adding scheduled message event (#30)
Closes #29

Reviewed-on: phoenix/textsender-api#30
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-11 16:43:51 +00:00
15 changed files with 1005 additions and 114 deletions
+7 -2
View File
@@ -69,7 +69,8 @@ func main() {
contactHandler := handler.NewContactHandler(contactStore) contactHandler := handler.NewContactHandler(contactStore)
messageHandler := handler.NewMessageHandler(messageStore) messageHandler := handler.NewMessageHandler(messageStore)
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore) scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore) scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore, schStore)
scheduledMessageStatusHandler := handler.NewScheduledMessageStatusHandler(schMsgEventStore, schStore)
router := chi.NewRouter() router := chi.NewRouter()
@@ -84,7 +85,11 @@ func main() {
router.Handle(endpoint.GET_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.GetMessage))) router.Handle(endpoint.GET_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.GetMessage)))
router.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage))) router.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage)))
router.Handle(endpoint.AddEventToScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.AddScheduledMessageEvent))) router.Handle(endpoint.AddEventToScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.AddScheduledMessageEvent)))
router.Handle(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("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 // Start server
server := &http.Server{ server := &http.Server{
+7 -3
View File
@@ -62,7 +62,8 @@ func TestMain(m *testing.M) {
contactHandler := handler.NewContactHandler(contactStore) contactHandler := handler.NewContactHandler(contactStore)
messageHandler := handler.NewMessageHandler(messageStore) messageHandler := handler.NewMessageHandler(messageStore)
scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore) scheduledMessageHandler := handler.NewScheduledMessageHandler(schStore)
scheduledMessageEventHandler := handler.NewScheduledMessageEventHandler(schMsgEventStore) 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)))
@@ -70,8 +71,11 @@ func TestMain(m *testing.M) {
testRouter.Handle(endpoint.ADD_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.AddMessage))) testRouter.Handle(endpoint.ADD_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.AddMessage)))
testRouter.Handle(endpoint.GET_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.GetMessage))) testRouter.Handle(endpoint.GET_MESSAGE, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(messageHandler.GetMessage)))
testRouter.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage))) testRouter.Handle(endpoint.ScheduleMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageHandler.AddScheduledMessage)))
testRouter.Handle(endpoint.AddEventToScheduledMessageEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.AddScheduledMessageEvent))) testRouter.Method("GET", endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent)))
testRouter.Handle(endpoint.GetScheduledMessageEventEndpoint, mdlware.AuthMiddleware(jwtService)(http.HandlerFunc(scheduledMessageEventHandler.GetScheduledMessageEvent))) 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() code := m.Run()
os.Exit(code) os.Exit(code)
+5 -1
View File
@@ -6,5 +6,9 @@ const GET_MESSAGE = "/api/v1/message"
const GET_CONTACT = "/api/v1/contact" const GET_CONTACT = "/api/v1/contact"
const ADD_CONTACT_ENDPOINT = "/api/v1/contact/new" 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 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 UpdateScheduledMessageStatusEndpoint = "/api/v1/schedule/message/status/update"
const FetchNextScheduledMessageEndpoint = "/api/v1/schedule/message/fetch"
+21 -3
View File
@@ -1,8 +1,11 @@
package handler package handler
import "encoding/json" import (
import "log" "encoding/json"
import "net/http" "fmt"
"log"
"net/http"
)
func ExtractFromRequest(r *http.Request, reqItem interface{}) error { func ExtractFromRequest(r *http.Request, reqItem interface{}) error {
err := json.NewDecoder(r.Body).Decode(&reqItem) err := json.NewDecoder(r.Body).Decode(&reqItem)
@@ -21,3 +24,18 @@ func RespondWithJSON(w http.ResponseWriter, statusCode int, data interface{}) {
log.Printf("Error encoding JSON: %v", err) log.Printf("Error encoding JSON: %v", err)
} }
} }
// Gets the query parameter from the URL
func ParseQueryParams(r *http.Request, query string) (*string, error) {
queryParams := r.URL.Query()
if _, exists := queryParams[query]; exists {
value := queryParams.Get(query)
if len(value) == 0 {
return nil, fmt.Errorf("Value of query parameter is empty")
} else {
return &value, nil
}
} else {
return nil, fmt.Errorf("Could not find query %s", query)
}
}
+154
View File
@@ -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")
}
+100 -4
View File
@@ -24,6 +24,16 @@ type AddScheduledMessageResponse struct {
Data []scheduling.ScheduledMessage `json:"data"` Data []scheduling.ScheduledMessage `json:"data"`
} }
type GetScheduledMessageResponse struct {
Message string `json:"message"`
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
} }
@@ -32,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
@@ -55,10 +65,10 @@ 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 = 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"
@@ -83,6 +93,92 @@ func (c *ScheduledMessageHandler) AddScheduledMessage(w http.ResponseWriter, r *
RespondWithJSON(w, statusCode, &resp) RespondWithJSON(w, statusCode, &resp)
} }
func (s *ScheduledMessageHandler) GetScheduledMessage(w http.ResponseWriter, r *http.Request) {
var id, userId uuid.UUID
if idParam, err := ParseQueryParams(r, "id"); err == nil {
var err error
id, err = uuid.Parse(*idParam)
if err != nil {
http.Error(w, "Error parsing Id", http.StatusBadRequest)
return
}
}
if userIdParam, err := ParseQueryParams(r, "user_id"); err == nil {
userId, err = uuid.Parse(*userIdParam)
if err != nil {
http.Error(w, "Error parsing Id", http.StatusBadRequest)
return
}
}
var resp GetScheduledMessageResponse
var statusCode int
ctx := r.Context()
if userId == uuid.Nil && id == uuid.Nil {
http.Error(w, "Invalid query parameters", http.StatusBadRequest)
return
} else if userId != uuid.Nil && id != uuid.Nil {
if schMsgs, err := s.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil {
statusCode = http.StatusInternalServerError
resp.Message = err.Error()
} else {
for _, schMsg := range schMsgs {
if schMsg.Id == id {
resp.Data = append(resp.Data, *schMsg)
}
}
statusCode = http.StatusOK
resp.Message = "Successful"
}
} else if id != uuid.Nil {
if schMsg, err := s.ScheduledMessageStore.Get(ctx, id); err != nil {
statusCode = http.StatusInternalServerError
resp.Message = err.Error()
} else {
statusCode = http.StatusOK
resp.Message = "Successful"
resp.Data = append(resp.Data, *schMsg)
}
} else {
if schMsgs, err := s.ScheduledMessageStore.GetWithUserId(ctx, userId); err != nil {
statusCode = http.StatusInternalServerError
resp.Message = err.Error()
} else {
statusCode = http.StatusOK
resp.Message = "Successful"
for _, schMsg := range schMsgs {
resp.Data = append(resp.Data, *schMsg)
}
}
}
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)
@@ -99,7 +195,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
+88 -6
View File
@@ -27,12 +27,18 @@ type GetScheduledMessageEventResponse struct {
Data []scheduling.ScheduledMessageEvent `json:"data"` Data []scheduling.ScheduledMessageEvent `json:"data"`
} }
type ScheduledMessageEventHandler struct { type DeleteScheduledMessageEventResponse struct {
ScheduledMessageEventStore store.ScheduledMessageEventStore Message string `json:"message"`
Data []scheduling.ScheduledMessageEvent `json:"data"`
} }
func NewScheduledMessageEventHandler(str store.ScheduledMessageEventStore) *ScheduledMessageEventHandler { type ScheduledMessageEventHandler struct {
return &ScheduledMessageEventHandler{ScheduledMessageEventStore: str} ScheduledMessageEventStore store.ScheduledMessageEventStore
ScheduledMessageStore store.ScheduledMessageStore
}
func NewScheduledMessageEventHandler(str store.ScheduledMessageEventStore, schStore store.ScheduledMessageStore) *ScheduledMessageEventHandler {
return &ScheduledMessageEventHandler{ScheduledMessageEventStore: str, ScheduledMessageStore: schStore}
} }
func (s *ScheduledMessageEventHandler) AddScheduledMessageEvent(w http.ResponseWriter, r *http.Request) { func (s *ScheduledMessageEventHandler) AddScheduledMessageEvent(w http.ResponseWriter, r *http.Request) {
@@ -64,6 +70,11 @@ func (s *ScheduledMessageEventHandler) AddScheduledMessageEvent(w http.ResponseW
} else { } else {
ctx := r.Context() ctx := r.Context()
if schMsg, err := s.ScheduledMessageStore.Get(ctx, event.ScheduledMessageId); err != nil {
statusCode = http.StatusInternalServerError
resp.Message = err.Error()
} else {
if schMsg.Status == scheduling.Pending {
if exists, err := s.ScheduledMessageEventStore.Exists(ctx, &event); err == nil { if exists, err := s.ScheduledMessageEventStore.Exists(ctx, &event); err == nil {
if exists { if exists {
statusCode = http.StatusBadRequest statusCode = http.StatusBadRequest
@@ -82,6 +93,11 @@ func (s *ScheduledMessageEventHandler) AddScheduledMessageEvent(w http.ResponseW
statusCode = http.StatusInternalServerError statusCode = http.StatusInternalServerError
resp.Message = err.Error() resp.Message = err.Error()
} }
} else {
statusCode = http.StatusBadRequest
resp.Message = "Scheduled message cannot be modified due to the status"
}
}
} }
RespondWithJSON(w, statusCode, &resp) RespondWithJSON(w, statusCode, &resp)
@@ -93,6 +109,57 @@ func (s *ScheduledMessageEventHandler) GetScheduledMessageEvent(w http.ResponseW
return return
} }
var id, scheduledMessageId uuid.UUID
if idParam, err := ParseQueryParams(r, "id"); err == nil {
id, err = uuid.Parse(*idParam)
}
if scheduledMessageIdParam, err := ParseQueryParams(r, "scheduled_message_id"); err == nil {
scheduledMessageId, err = uuid.Parse(*scheduledMessageIdParam)
}
var statusCode int
var resp GetScheduledMessageEventResponse
ctx := r.Context()
if id == uuid.Nil && scheduledMessageId == uuid.Nil {
statusCode = http.StatusBadRequest
resp.Message = "Query parameters missing"
} else if id != uuid.Nil {
if event, err := s.ScheduledMessageEventStore.Get(ctx, id); err != nil {
resp.Message = err.Error()
statusCode = http.StatusInternalServerError
} else {
if event != nil {
resp.Message = "Successful"
statusCode = http.StatusOK
resp.Data = append(resp.Data, *event)
} else {
statusCode = http.StatusNotFound
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)
}
func (s *ScheduledMessageEventHandler) DeleteScheduledMessageEvent(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodDelete {
http.Error(w, "Metnot allowed", http.StatusMethodNotAllowed)
return
}
id := chi.URLParam(r, "id") id := chi.URLParam(r, "id")
if len(id) == 0 { if len(id) == 0 {
pathParts := strings.Split(r.URL.Path, "/") pathParts := strings.Split(r.URL.Path, "/")
@@ -104,8 +171,8 @@ func (s *ScheduledMessageEventHandler) GetScheduledMessageEvent(w http.ResponseW
} }
} }
var resp DeleteScheduledMessageEventResponse
var statusCode int var statusCode int
var resp GetScheduledMessageEventResponse
if parsedId, err := uuid.Parse(id); err != nil { if parsedId, err := uuid.Parse(id); err != nil {
resp.Message = err.Error() resp.Message = err.Error()
@@ -115,10 +182,25 @@ func (s *ScheduledMessageEventHandler) GetScheduledMessageEvent(w http.ResponseW
if event, err := s.ScheduledMessageEventStore.Get(ctx, parsedId); err != nil { if event, err := s.ScheduledMessageEventStore.Get(ctx, parsedId); err != nil {
resp.Message = err.Error() resp.Message = err.Error()
statusCode = http.StatusInternalServerError statusCode = http.StatusInternalServerError
} else {
if schMsg, err := s.ScheduledMessageStore.Get(ctx, event.ScheduledMessageId); err != nil {
resp.Message = err.Error()
statusCode = http.StatusInternalServerError
} else {
if schMsg.Status == scheduling.Pending {
if err = s.ScheduledMessageEventStore.Delete(ctx, event.Id); err != nil {
resp.Message = err.Error()
statusCode = http.StatusInternalServerError
} else { } else {
resp.Message = "Successful" resp.Message = "Successful"
statusCode = http.StatusOK
resp.Data = append(resp.Data, *event) resp.Data = append(resp.Data, *event)
statusCode = http.StatusOK
}
} else {
statusCode = http.StatusBadRequest
resp.Message = "Invalid status"
}
}
} }
} }
@@ -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"
@@ -18,6 +19,13 @@ import (
"git.kundeng.us/phoenix/textsender-api/internal/store/mock" "git.kundeng.us/phoenix/textsender-api/internal/store/mock"
) )
var (
recipientId = uuid.New()
messageId = uuid.New()
scheduledMessageId = uuid.New()
testUserId = uuid.New()
)
type CreateScheduledMessageEventRequest struct { type CreateScheduledMessageEventRequest struct {
RecipientId uuid.UUID `json:"recipient_id"` RecipientId uuid.UUID `json:"recipient_id"`
MessageId uuid.UUID `json:"message_id"` MessageId uuid.UUID `json:"message_id"`
@@ -32,27 +40,16 @@ func TestCreateScheduledMessageEventWithMock(t *testing.T) {
messageStore := mock.NewMockMessageStore() messageStore := mock.NewMockMessageStore()
schMsgStore := mock.NewMockScheduledMessageStore() schMsgStore := mock.NewMockScheduledMessageStore()
handler := NewScheduledMessageEventHandler(mockStore) handler := NewScheduledMessageEventHandler(mockStore, schMsgStore)
recipientId := uuid.New() recipientId := uuid.New()
messageId := uuid.New() messageId := uuid.New()
scheduledMessageId := uuid.New() scheduledMessageId := uuid.New()
testUserId := uuid.New() testUserId := uuid.New()
con := contact.Contact{} con := testContact(recipientId, testUserId)
con.Id = recipientId msg := testMessage(messageId, testUserId)
con.PhoneNumber = "+10123456789" schMsg := testScheduledMessage(scheduledMessageId, testUserId, now)
con.UserId = testUserId
msg := message.Message{}
msg.Id = messageId
msg.Content = "Oh how the might have fallen"
msg.UserId = testUserId
schMsg := scheduling.ScheduledMessage{}
schMsg.Id = scheduledMessageId
schMsg.UserId = testUserId
schMsg.Scheduled = now.Add(20 * time.Minute)
ctx := t.Context() ctx := t.Context()
@@ -93,32 +90,17 @@ func TestGetScheduledMessageEventWithMock(t *testing.T) {
messageStore := mock.NewMockMessageStore() messageStore := mock.NewMockMessageStore()
schMsgStore := mock.NewMockScheduledMessageStore() schMsgStore := mock.NewMockScheduledMessageStore()
schMsgEventStore := mock.NewMockScheduledMessageEventStore() schMsgEventStore := mock.NewMockScheduledMessageEventStore()
handler := NewScheduledMessageEventHandler(schMsgEventStore) handler := NewScheduledMessageEventHandler(schMsgEventStore, schMsgStore)
recipientId := uuid.New() recipientId := uuid.New()
messageId := uuid.New() messageId := uuid.New()
scheduledMessageId := uuid.New() scheduledMessageId := uuid.New()
testUserId := uuid.New() testUserId := uuid.New()
con := contact.Contact{} con := testContact(recipientId, testUserId)
con.Id = recipientId msg := testMessage(messageId, testUserId)
con.PhoneNumber = "+10123456789" schMsg := testScheduledMessage(scheduledMessageId, testUserId, now)
con.UserId = testUserId event := testScheduledMessageEvent(msg.Id, con.Id, schMsg.Id)
msg := message.Message{}
msg.Id = messageId
msg.Content = "Oh how the might have fallen"
msg.UserId = testUserId
schMsg := scheduling.ScheduledMessage{}
schMsg.Id = scheduledMessageId
schMsg.UserId = testUserId
schMsg.Scheduled = now.Add(20 * time.Minute)
event := scheduling.ScheduledMessageEvent{}
event.MessageId = msg.Id
event.RecipientId = con.Id
event.ScheduledMessageId = schMsg.Id
ctx := t.Context() ctx := t.Context()
@@ -132,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)
@@ -150,3 +132,75 @@ func TestGetScheduledMessageEventWithMock(t *testing.T) {
assert.NotEmpty(t, msgEvent.Created, "Created date should not be empty") assert.NotEmpty(t, msgEvent.Created, "Created date should not be empty")
} }
func TestDeleteScheduledMessageEventWithMock(t *testing.T) {
now := time.Now()
contactStore := mock.NewMockContactStore()
messageStore := mock.NewMockMessageStore()
schMsgStore := mock.NewMockScheduledMessageStore()
schMsgEventStore := mock.NewMockScheduledMessageEventStore()
handler := NewScheduledMessageEventHandler(schMsgEventStore, schMsgStore)
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)
}
endpointValue := strings.Replace(endpoint.DeleteScheduledMessageEventEndpoint, "{id}", event.Id.String(), 1)
req, _ := http.NewRequest("DELETE", endpointValue, nil)
rr := httptest.NewRecorder()
handler.DeleteScheduledMessageEvent(rr, req)
assert.Equal(t, http.StatusOK, rr.Code)
var response DeleteScheduledMessageEventResponse
err := json.Unmarshal(rr.Body.Bytes(), &response)
assert.NoError(t, err, "Error creating event %v", err)
assert.NotEmpty(t, response.Data, "No event created")
msgEvent := response.Data[0]
assert.NotEmpty(t, msgEvent.Created, "Created date should not be empty")
}
func testContact(id uuid.UUID, userId uuid.UUID) contact.Contact {
if id == uuid.Nil {
return contact.Contact{Id: uuid.New(), PhoneNumber: "+10123456789", UserId: userId}
} else {
return contact.Contact{Id: id, PhoneNumber: "+10123456789", UserId: userId}
}
}
func testMessage(id uuid.UUID, userId uuid.UUID) message.Message {
if id == uuid.Nil {
return message.Message{Id: uuid.New(), Content: "Oh how the mighty have fallen", UserId: userId}
} else {
return message.Message{Id: id, Content: "Oh how the mighty have fallen", UserId: userId}
}
}
func testScheduledMessage(id uuid.UUID, userId uuid.UUID, now time.Time) scheduling.ScheduledMessage {
if id == uuid.Nil {
return scheduling.ScheduledMessage{Id: uuid.New(), UserId: userId, Scheduled: now.Add(20 * time.Minute), Status: scheduling.Pending}
} else {
return scheduling.ScheduledMessage{Id: id, UserId: userId, Scheduled: now.Add(20 * time.Minute), Status: scheduling.Pending}
}
}
func testScheduledMessageEvent(messageId, recipientId, scheduledMessageId uuid.UUID) scheduling.ScheduledMessageEvent {
return scheduling.ScheduledMessageEvent{MessageId: messageId, RecipientId: recipientId, ScheduledMessageId: scheduledMessageId}
}
+88 -2
View File
@@ -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"
@@ -28,8 +29,7 @@ func TestCreateScheduledMessageWithMock(t *testing.T) {
handler := NewScheduledMessageHandler(mockStore) handler := NewScheduledMessageHandler(mockStore)
testUserId := uuid.New() testUserId := uuid.New()
scheduled := now.Add(5 * time.Minute) testBody := testCreateScheduledMessageRequest(testUserId, now)
testBody := CreateScheduledMessageRequest{Status: scheduling.Pending, UserId: testUserId, Scheduled: scheduled}
jsonValue, _ := json.Marshal(testBody) jsonValue, _ := json.Marshal(testBody)
req, _ := http.NewRequest("POST", endpoint.ScheduleMessageEndpoint, strings.NewReader(string(jsonValue))) req, _ := http.NewRequest("POST", endpoint.ScheduleMessageEndpoint, strings.NewReader(string(jsonValue)))
@@ -47,3 +47,89 @@ func TestCreateScheduledMessageWithMock(t *testing.T) {
assert.NotNil(t, response.Data[0].Id, "Id should not be nil") assert.NotNil(t, response.Data[0].Id, "Id should not be nil")
} }
func TestGetScheduledMessageWithMock(t *testing.T) {
now := time.Now()
mockStore := mock.NewMockScheduledMessageStore()
testUserId := uuid.New()
testBody := testCreateScheduledMessageRequest(testUserId, now)
ctx := t.Context()
schMsg := scheduling.ScheduledMessage{}
schMsg.Scheduled = testBody.Scheduled
schMsg.Status = testBody.Status
schMsg.UserId = testBody.UserId
if err := mockStore.CreateScheduledMessage(ctx, &schMsg); err != nil {
assert.NoError(t, err, "Error Creating message %v", err)
}
url := fmt.Sprintf("%s?id=%s&user_id=%s", endpoint.GetScheduledMessageEndpoint, schMsg.Id, schMsg.UserId)
req, _ := http.NewRequest("GET", url, nil)
rr := httptest.NewRecorder()
handler := NewScheduledMessageHandler(mockStore)
handler.GetScheduledMessage(rr, req)
assert.Equal(t, http.StatusOK, rr.Code)
var response GetScheduledMessageResponse
err := json.Unmarshal(rr.Body.Bytes(), &response)
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}
}
-38
View File
@@ -9,7 +9,6 @@ import (
"git.kundeng.us/phoenix/textsender-models/pkg/contact" "git.kundeng.us/phoenix/textsender-models/pkg/contact"
"git.kundeng.us/phoenix/textsender-models/pkg/message" "git.kundeng.us/phoenix/textsender-models/pkg/message"
"git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling"
) )
type Key struct { type Key struct {
@@ -201,40 +200,3 @@ func (m *MockMessageStore) MessageExists(ctx context.Context, msg *message.Messa
return exists, nil return exists, nil
} }
type ScheduledMessageKey struct {
UserId uuid.UUID
}
type MockScheduledMessageStore struct {
ScheduledMessages map[uuid.UUID]*scheduling.ScheduledMessage
ScheduledMessagesByKey map[ScheduledMessageKey]*scheduling.ScheduledMessage
mu sync.RWMutex
Error error // Optional: simulate errors
}
func NewMockScheduledMessageStore() *MockScheduledMessageStore {
return &MockScheduledMessageStore{
ScheduledMessages: make(map[uuid.UUID]*scheduling.ScheduledMessage),
ScheduledMessagesByKey: make(map[ScheduledMessageKey]*scheduling.ScheduledMessage),
}
}
func (m *MockScheduledMessageStore) CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.Error != nil {
return m.Error
}
if schedMsg.Id == uuid.Nil {
schedMsg.Id = uuid.New()
}
key := ScheduledMessageKey{UserId: schedMsg.UserId}
m.ScheduledMessages[schedMsg.Id] = schedMsg
m.ScheduledMessagesByKey[key] = schedMsg
return nil
}
@@ -46,11 +46,81 @@ func (m *MockScheduledMessageEventStore) Get(ctx context.Context, id uuid.UUID)
if _, exists := m.ScheduledMessageEvents[id]; exists { if _, exists := m.ScheduledMessageEvents[id]; exists {
return m.ScheduledMessageEvents[id], nil return m.ScheduledMessageEvents[id], nil
} else { } else {
fmt.Println("Not found")
return nil, fmt.Errorf("Not found") return nil, fmt.Errorf("Not found")
} }
} }
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 {
m.mu.Lock()
defer m.mu.Unlock()
if m.Error != nil {
return m.Error
}
if id == uuid.Nil {
return fmt.Errorf("Id is nil")
}
if _, exists := m.ScheduledMessageEvents[id]; exists {
originalAmount := len(m.ScheduledMessageEvents)
copiedEvents := make(map[uuid.UUID]*scheduling.ScheduledMessageEvent)
copiedEventsKey := make(map[ScheduledMessageEventKey]*scheduling.ScheduledMessageEvent)
for i, schMsgEvent := range m.ScheduledMessageEvents {
key := ScheduledMessageEventKey{RecipientId: schMsgEvent.RecipientId, MessageId: schMsgEvent.MessageId, ScheduledMessageId: schMsgEvent.ScheduledMessageId}
if schMsgEvent.Id != id {
copiedEvents[i] = schMsgEvent
copiedEventsKey[key] = schMsgEvent
}
}
if originalAmount > 1 {
if len(copiedEvents) > 0 && len(copiedEventsKey) > 0 {
m.ScheduledMessageEvents = copiedEvents
m.ScheduledMessageEventsByKey = copiedEventsKey
return nil
} else {
return fmt.Errorf("Not removed")
}
} else {
m.ScheduledMessageEvents = make(map[uuid.UUID]*scheduling.ScheduledMessageEvent)
m.ScheduledMessageEventsByKey = make(map[ScheduledMessageEventKey]*scheduling.ScheduledMessageEvent)
return nil
}
} else {
return fmt.Errorf("Not found")
}
}
func (m *MockScheduledMessageEventStore) CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error { func (m *MockScheduledMessageEventStore) CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() defer m.mu.Unlock()
@@ -0,0 +1,142 @@
package mock
import (
"context"
"fmt"
"sync"
"github.com/google/uuid"
"git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling"
)
type ScheduledMessageKey struct {
UserId uuid.UUID
}
type MockScheduledMessageStore struct {
ScheduledMessages map[uuid.UUID]*scheduling.ScheduledMessage
ScheduledMessagesByKey map[ScheduledMessageKey]*scheduling.ScheduledMessage
mu sync.RWMutex
Error error // Optional: simulate errors
}
func NewMockScheduledMessageStore() *MockScheduledMessageStore {
return &MockScheduledMessageStore{
ScheduledMessages: make(map[uuid.UUID]*scheduling.ScheduledMessage),
ScheduledMessagesByKey: make(map[ScheduledMessageKey]*scheduling.ScheduledMessage),
}
}
func (m *MockScheduledMessageStore) Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessage, 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 {
return schMsg, nil
} else {
return nil, fmt.Errorf("Scheduled message does not exist")
}
}
}
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()
if m.Error != nil {
return nil, m.Error
}
if userId == uuid.Nil {
return nil, fmt.Errorf("User Id is nil")
} else {
var schMsgs []*scheduling.ScheduledMessage
for _, schMsg := range m.ScheduledMessages {
if schMsg.UserId == userId {
schMsgs = append(schMsgs, schMsg)
}
}
if len(schMsgs) == 0 {
return nil, nil
} else {
return schMsgs, nil
}
}
}
func (m *MockScheduledMessageStore) CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.Error != nil {
return m.Error
}
if schedMsg.Id == uuid.Nil {
schedMsg.Id = uuid.New()
}
key := ScheduledMessageKey{UserId: schedMsg.UserId}
m.ScheduledMessages[schedMsg.Id] = schedMsg
m.ScheduledMessagesByKey[key] = schedMsg
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,8 @@ 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
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)
} }
@@ -35,12 +37,56 @@ func (s *PGScheduledMessageEventStore) Get(ctx context.Context, id uuid.UUID) (*
if err == pgx.ErrNoRows { if err == pgx.ErrNoRows {
return nil, nil return nil, nil
} } else if err != nil {
if err != nil {
return nil, fmt.Errorf("getting scheduled message event by ID: %w", err) return nil, fmt.Errorf("getting scheduled message event by ID: %w", err)
} else {
return &event, nil
}
}
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)
} }
return &event, nil 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 {
query := `
DELETE FROM scheduled_message_events WHERE id = $1
`
commandTag, err := s.db.Exec(ctx, query, id)
if err != nil {
return fmt.Errorf("error deleting event: %w", err)
}
if commandTag.RowsAffected() == 0 {
return fmt.Errorf("no event found with id %d", id)
} else {
return nil
}
} }
func (s *PGScheduledMessageEventStore) CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error { func (s *PGScheduledMessageEventStore) CreateScheduledMessageEvent(ctx context.Context, event *scheduling.ScheduledMessageEvent) error {
+95
View File
@@ -2,14 +2,21 @@ package store
import ( import (
"context" "context"
"fmt"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool" "github.com/jackc/pgx/v5/pgxpool"
"git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling" "git.kundeng.us/phoenix/textsender-models/pkg/message/scheduling"
) )
type ScheduledMessageStore interface { 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 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 {
@@ -20,6 +27,79 @@ func NewScheduledMessageStore(db *pgxpool.Pool) *PGScheduledMessageStore {
return &PGScheduledMessageStore{db: db} return &PGScheduledMessageStore{db: db}
} }
func (s *PGScheduledMessageStore) Get(ctx context.Context, id uuid.UUID) (*scheduling.ScheduledMessage, error) {
query := `
SELECT id, scheduled, created, status, user_id FROM scheduled_messages WHERE id = $1
`
var schMsg scheduling.ScheduledMessage
err := s.db.QueryRow(ctx, query, id).Scan(
&schMsg.Id, &schMsg.Scheduled, &schMsg.Created, &schMsg.Status, &schMsg.UserId,
)
if err == pgx.ErrNoRows {
return nil, fmt.Errorf("No rows")
} else if err != nil {
return nil, fmt.Errorf("Getting scheduled message: %w", err)
}
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
`
rows, err := s.db.Query(ctx, query, userId)
if err != nil {
return nil, fmt.Errorf("Error retrieving rows from scheduled_messages: %w", err)
}
var schMsgs []*scheduling.ScheduledMessage
for rows.Next() {
var schMsg scheduling.ScheduledMessage
if err := rows.Scan(&schMsg.Id, &schMsg.Scheduled, &schMsg.Created, &schMsg.Status, &schMsg.UserId); err != nil {
return nil, fmt.Errorf("Error fetching data from row: %w", err)
} else {
schMsgs = append(schMsgs, &schMsg)
}
}
if err := rows.Err(); err != nil {
return nil, err
} else {
return schMsgs, nil
}
}
func (s *PGScheduledMessageStore) CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error { func (s *PGScheduledMessageStore) CreateScheduledMessage(ctx context.Context, schedMsg *scheduling.ScheduledMessage) error {
query := ` query := `
INSERT INTO scheduled_messages (scheduled, status, user_id) INSERT INTO scheduled_messages (scheduled, status, user_id)
@@ -31,3 +111,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
}
}