Refactor slack messages

Signed-off-by: Sacha Al Himdani <sacha@getprobo.com>
This commit is contained in:
Sacha Al Himdani
2025-10-23 09:32:55 +02:00
parent 0073846d3b
commit b0f60d8de5
9 changed files with 611 additions and 529 deletions

View File

@@ -81,7 +81,8 @@ func (s *Sender) batchSendMessages(ctx context.Context) error {
message.Error = &panicErr
message.UpdatedAt = time.Now()
if updateErr := message.Update(ctx, tx); updateErr != nil {
scope := coredata.NewScope(message.ID.TenantID())
if updateErr := message.Update(ctx, tx, scope); updateErr != nil {
s.logger.ErrorCtx(ctx, "cannot update slack message after panic", log.Error(updateErr))
}
@@ -90,24 +91,32 @@ func (s *Sender) batchSendMessages(ctx context.Context) error {
}
}()
err = message.LoadNextUnsentForUpdate(ctx, tx)
err = message.LoadNextInitalUnsentForUpdate(ctx, tx)
if err != nil {
return err
}
scope := coredata.NewScope(message.ID.TenantID())
channelID, messageTS, sendErr := s.sendMessage(ctx, tx, message)
message.ChannelID = channelID
message.MessageTS = messageTS
now := time.Now()
message.UpdatedAt = now
if channelID != nil && messageTS != nil {
message.ChannelID = channelID
message.MessageTS = messageTS
message.UpdatedAt = now
if err := message.UpdateChannelAndTSByInitialMessageID(ctx, tx, scope, message.ID, *channelID, *messageTS, now); err != nil {
return fmt.Errorf("cannot update all messages with initial message id: %w", err)
}
}
if sendErr != nil {
errorMsg := sendErr.Error()
message.Error = &errorMsg
message.UpdatedAt = time.Now()
if err := message.Update(ctx, tx); err != nil {
if err := message.Update(ctx, tx, scope); err != nil {
return fmt.Errorf("cannot update slack message with error: %w", err)
}
@@ -117,7 +126,7 @@ func (s *Sender) batchSendMessages(ctx context.Context) error {
message.SentAt = &now
if err := message.Update(ctx, tx); err != nil {
if err := message.Update(ctx, tx, scope); err != nil {
return fmt.Errorf("cannot update slack message: %w", err)
}
@@ -196,62 +205,59 @@ func (s *Sender) batchUpdateMessages(ctx context.Context) error {
err := s.pg.WithTx(
ctx,
func(tx pg.Conn) (err error) {
update := &coredata.SlackMessageUpdate{}
updateMessage := &coredata.SlackMessage{}
defer func() {
if r := recover(); r != nil {
panicErr := fmt.Sprintf("panic recovered: %v", r)
update.Error = &panicErr
update.UpdatedAt = time.Now()
updateMessage.Error = &panicErr
updateMessage.UpdatedAt = time.Now()
if updateErr := update.Update(ctx, tx); updateErr != nil {
s.logger.ErrorCtx(ctx, "cannot update slack message update after panic", log.Error(updateErr))
scope := coredata.NewScope(updateMessage.ID.TenantID())
if updateErr := updateMessage.Update(ctx, tx, scope); updateErr != nil {
s.logger.ErrorCtx(ctx, "cannot update slack message after panic", log.Error(updateErr))
}
s.logger.ErrorCtx(ctx, "panic while updating slack message", log.String("error", panicErr), log.String("update_id", update.ID.String()))
s.logger.ErrorCtx(ctx, "panic while updating slack message", log.String("error", panicErr), log.String("message_id", updateMessage.ID.String()))
err = fmt.Errorf("panic recovered: %v", r)
}
}()
err = update.LoadNextUnsentForUpdate(ctx, tx)
err = updateMessage.LoadNextUpdateUnsentForUpdate(ctx, tx)
if err != nil {
return err
}
message := &coredata.SlackMessage{}
if err := message.LoadById(ctx, tx, coredata.NewScope(update.SlackMessageID.TenantID()), update.SlackMessageID); err != nil {
return fmt.Errorf("cannot load slack message: %w", err)
}
updateErr := s.updateMessage(ctx, tx, message, update)
scope := coredata.NewScope(updateMessage.ID.TenantID())
updateErr := s.updateMessage(ctx, tx, updateMessage)
now := time.Now()
update.UpdatedAt = now
updateMessage.UpdatedAt = now
if updateErr != nil {
errorMsg := updateErr.Error()
update.Error = &errorMsg
update.UpdatedAt = time.Now()
updateMessage.Error = &errorMsg
updateMessage.UpdatedAt = time.Now()
if err := update.Update(ctx, tx); err != nil {
return fmt.Errorf("cannot update slack message update with error: %w", err)
if err := updateMessage.Update(ctx, tx, scope); err != nil {
return fmt.Errorf("cannot update slack message with error: %w", err)
}
s.logger.ErrorCtx(ctx, "error updating slack message", log.Error(updateErr), log.String("update_id", update.ID.String()))
s.logger.ErrorCtx(ctx, "error updating slack message", log.Error(updateErr), log.String("message_id", updateMessage.ID.String()))
return nil
}
update.SentAt = &now
updateMessage.SentAt = &now
if err := update.Update(ctx, tx); err != nil {
return fmt.Errorf("cannot update slack message update: %w", err)
if err := updateMessage.Update(ctx, tx, scope); err != nil {
return fmt.Errorf("cannot update slack message: %w", err)
}
return nil
},
)
if errors.Is(err, coredata.ErrNoUnsentSlackMessageUpdate{}) {
if errors.Is(err, coredata.ErrNoUnsentSlackMessage{}) {
return nil
}
@@ -261,12 +267,12 @@ func (s *Sender) batchUpdateMessages(ctx context.Context) error {
}
}
func (s *Sender) updateMessage(ctx context.Context, tx pg.Conn, message *coredata.SlackMessage, update *coredata.SlackMessageUpdate) error {
if message.ChannelID == nil || message.MessageTS == nil {
func (s *Sender) updateMessage(ctx context.Context, tx pg.Conn, updateMessage *coredata.SlackMessage) error {
if updateMessage.ChannelID == nil || updateMessage.MessageTS == nil {
return fmt.Errorf("slack message has no channel ID or message TS")
}
tenantID := message.ID.TenantID()
tenantID := updateMessage.ID.TenantID()
scope := coredata.NewScope(tenantID)
var connectors coredata.Connectors
@@ -274,7 +280,7 @@ func (s *Sender) updateMessage(ctx context.Context, tx pg.Conn, message *coredat
ctx,
tx,
scope,
message.OrganizationID,
updateMessage.OrganizationID,
coredata.ConnectorProtocolOAuth2,
coredata.ConnectorProviderSlack,
s.encryptionKey,
@@ -302,7 +308,7 @@ func (s *Sender) updateMessage(ctx context.Context, tx pg.Conn, message *coredat
client := NewClient(s.logger)
if err := client.UpdateMessage(ctx, slackConn.AccessToken, *message.ChannelID, *message.MessageTS, update.Body); err != nil {
if err := client.UpdateMessage(ctx, slackConn.AccessToken, *updateMessage.ChannelID, *updateMessage.MessageTS, updateMessage.Body); err != nil {
s.logger.ErrorCtx(ctx, "failed to update message on Slack", log.Error(err))
return fmt.Errorf("failed to update message on Slack: %w", err)
}