feat(webhooks): align account payloads

This commit is contained in:
2026-06-06 06:26:30 +08:00
parent 1adf1c9c31
commit 02ea135e0c
13 changed files with 495 additions and 94 deletions
@@ -1,10 +1,14 @@
package v1
import (
"encoding/json"
"errors"
"net/http"
"github.com/gin-gonic/gin"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/service"
"github.com/gochat/gochat/pkg/response"
)
@@ -35,7 +39,7 @@ func (h *WebhookSubscriptionHandler) List(c *gin.Context) {
return
}
response.OK(c, gin.H{"webhook_subscriptions": subscriptions})
c.JSON(http.StatusOK, gin.H{"payload": gin.H{"webhooks": serializeWebhookSubscriptions(subscriptions)}})
}
// Get returns a single webhook subscription by ID.
@@ -45,19 +49,20 @@ func (h *WebhookSubscriptionHandler) Get(c *gin.Context) {
response.AbortWithStatusError(c, http.StatusServiceUnavailable, response.ErrInternal, "Webhook subscription service not available")
return
}
accountID := getAccountID(c)
webhookID, err := parseUintParam(c, "webhook_id")
if err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid webhook id")
return
}
subscription, svcErr := h.webhookSubscriptionService.GetSubscription(c.Request.Context(), webhookID)
subscription, svcErr := h.webhookSubscriptionService.GetWebhook(c.Request.Context(), accountID, webhookID)
if svcErr != nil {
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "Failed to fetch webhook subscription")
abortWebhookSubscriptionError(c, svcErr)
return
}
response.OK(c, gin.H{"webhook_subscription": subscription})
c.JSON(http.StatusOK, gin.H{"payload": gin.H{"webhook": serializeWebhookSubscription(*subscription)}})
}
// Create adds a new webhook subscription for an account.
@@ -69,22 +74,19 @@ func (h *WebhookSubscriptionHandler) Create(c *gin.Context) {
}
accountID := getAccountID(c)
var req struct {
URL string `json:"url" binding:"required"`
Events []string `json:"events" binding:"required"`
}
if err := c.ShouldBindJSON(&req); err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "Invalid request: url and events are required")
var req service.WebhookSubscriptionMutation
if err := bindJSONWrappedOrRaw(c, "webhook", &req); err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "Invalid request body")
return
}
subscription, err := h.webhookSubscriptionService.CreateSubscription(c.Request.Context(), accountID, req.URL, req.Events)
subscription, err := h.webhookSubscriptionService.CreateWebhook(c.Request.Context(), accountID, req)
if err != nil {
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "Failed to create webhook subscription")
abortWebhookSubscriptionError(c, err)
return
}
response.Created(c, subscription)
c.JSON(http.StatusOK, gin.H{"payload": gin.H{"webhook": serializeWebhookSubscription(*subscription)}})
}
// Update modifies a webhook subscription.
@@ -94,29 +96,26 @@ func (h *WebhookSubscriptionHandler) Update(c *gin.Context) {
response.AbortWithStatusError(c, http.StatusServiceUnavailable, response.ErrInternal, "Webhook subscription service not available")
return
}
id, err := parseUintParam(c, "id")
accountID := getAccountID(c)
id, err := parseUintAnyParam(c, "webhook_id", "id")
if err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "Invalid webhook subscription ID")
return
}
var req struct {
URL string `json:"url"`
Events []string `json:"events"`
Active bool `json:"active"`
}
if err := c.ShouldBindJSON(&req); err != nil {
var req service.WebhookSubscriptionMutation
if err := bindJSONWrappedOrRaw(c, "webhook", &req); err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "Invalid request body")
return
}
subscription, err := h.webhookSubscriptionService.UpdateSubscription(c.Request.Context(), id, req.URL, req.Events, req.Active)
subscription, err := h.webhookSubscriptionService.UpdateWebhook(c.Request.Context(), accountID, id, req)
if err != nil {
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "Failed to update webhook subscription")
abortWebhookSubscriptionError(c, err)
return
}
response.OK(c, subscription)
c.JSON(http.StatusOK, gin.H{"payload": gin.H{"webhook": serializeWebhookSubscription(*subscription)}})
}
// Delete removes a webhook subscription.
@@ -126,18 +125,19 @@ func (h *WebhookSubscriptionHandler) Delete(c *gin.Context) {
response.AbortWithStatusError(c, http.StatusServiceUnavailable, response.ErrInternal, "Webhook subscription service not available")
return
}
id, err := parseUintParam(c, "id")
accountID := getAccountID(c)
id, err := parseUintAnyParam(c, "webhook_id", "id")
if err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "Invalid webhook subscription ID")
return
}
if err := h.webhookSubscriptionService.DeleteSubscription(c.Request.Context(), id); err != nil {
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "Failed to delete webhook subscription")
if err := h.webhookSubscriptionService.DeleteWebhook(c.Request.Context(), accountID, id); err != nil {
abortWebhookSubscriptionError(c, err)
return
}
response.NoContent(c)
c.Status(http.StatusOK)
}
// ListDeliveries returns recent webhook delivery records for a subscription.
@@ -160,4 +160,41 @@ func (h *WebhookSubscriptionHandler) ListDeliveries(c *gin.Context) {
}
response.OK(c, gin.H{"deliveries": deliveries})
}
}
func abortWebhookSubscriptionError(c *gin.Context, err error) {
if errors.Is(err, gorm.ErrRecordNotFound) {
response.AbortWithStatusError(c, http.StatusNotFound, response.ErrNotFound, "webhook not found")
return
}
response.AbortWithStatusError(c, http.StatusUnprocessableEntity, response.ErrValidation, err.Error())
}
func serializeWebhookSubscriptions(subscriptions []model.WebhookSubscription) []gin.H {
items := make([]gin.H, 0, len(subscriptions))
for _, subscription := range subscriptions {
items = append(items, serializeWebhookSubscription(subscription))
}
return items
}
func serializeWebhookSubscription(subscription model.WebhookSubscription) gin.H {
var subscriptions []string
_ = json.Unmarshal(subscription.Events, &subscriptions)
payload := gin.H{
"id": subscription.ID,
"name": subscription.Name,
"url": subscription.URL,
"account_id": subscription.AccountID,
"subscriptions": subscriptions,
"secret": subscription.Secret,
}
if subscription.InboxID != nil && *subscription.InboxID != 0 {
inbox := gin.H{"id": *subscription.InboxID}
if subscription.Inbox.ID != 0 {
inbox["name"] = subscription.Inbox.Name
}
payload["inbox"] = inbox
}
return payload
}
@@ -2,6 +2,8 @@ package v1
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
@@ -31,7 +33,7 @@ func (s *WebhookSubscriptionHandlerTestSuite) SetupSuite() {
Logger: logger.Default.LogMode(logger.Silent),
})
s.Require().NoError(err)
s.Require().NoError(db.AutoMigrate(&model.Account{}, &model.WebhookSubscription{}))
s.Require().NoError(db.AutoMigrate(&model.Account{}, &model.Inbox{}, &model.WebhookSubscription{}))
s.db = db
repo := repository.NewWebhookSubscriptionRepo(db)
@@ -55,21 +57,98 @@ func TestWebhookSubscriptionHandlerSuite(t *testing.T) {
func (s *WebhookSubscriptionHandlerTestSuite) TestList_Success() {
r := gin.New()
r.GET("/api/v1/accounts/:account_id/webhooks/:webhook_id/subscriptions", s.handler.List)
r.GET("/api/v1/accounts/:account_id/webhooks", s.handler.List)
_, err := s.handler.webhookSubscriptionService.CreateWebhook(context.Background(), s.account.ID, service.WebhookSubscriptionMutation{
Name: "List hook",
URL: "https://example.com/list-hook",
Subscriptions: []string{"message_created"},
})
s.Require().NoError(err)
w := httptest.NewRecorder()
req, _ := http.NewRequest("GET", fmt.Sprintf("/api/v1/accounts/%d/webhooks/1/subscriptions", s.account.ID), nil)
req, _ := http.NewRequest("GET", fmt.Sprintf("/api/v1/accounts/%d/webhooks", s.account.ID), nil)
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusOK, w.Code)
var body map[string]any
s.Require().NoError(json.Unmarshal(w.Body.Bytes(), &body))
payload := body["payload"].(map[string]any)
webhooks := payload["webhooks"].([]any)
s.NotEmpty(webhooks)
}
func (s *WebhookSubscriptionHandlerTestSuite) TestCreate_Success_ChatwootPayload() {
r := gin.New()
r.POST("/api/v1/accounts/:account_id/webhooks", s.handler.Create)
body := `{"webhook":{"name":"Created hook","url":"https://example.com/created-hook","subscriptions":["conversation_created","message_created"]}}`
w := httptest.NewRecorder()
req, _ := http.NewRequest("POST", fmt.Sprintf("/api/v1/accounts/%d/webhooks", s.account.ID), bytes.NewBufferString(body))
req.Header.Set("Content-Type", "application/json")
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusOK, w.Code)
var parsed map[string]any
s.Require().NoError(json.Unmarshal(w.Body.Bytes(), &parsed))
webhook := parsed["payload"].(map[string]any)["webhook"].(map[string]any)
s.Equal("Created hook", webhook["name"])
s.Equal("https://example.com/created-hook", webhook["url"])
s.NotEmpty(webhook["secret"])
s.Equal([]any{"conversation_created", "message_created"}, webhook["subscriptions"])
}
func (s *WebhookSubscriptionHandlerTestSuite) TestUpdate_Success_ChatwootPayload() {
created, err := s.handler.webhookSubscriptionService.CreateWebhook(context.Background(), s.account.ID, service.WebhookSubscriptionMutation{
Name: "Before",
URL: "https://example.com/update-before",
Subscriptions: []string{"message_created"},
})
s.Require().NoError(err)
r := gin.New()
r.PATCH("/api/v1/accounts/:account_id/webhooks/:webhook_id", s.handler.Update)
body := `{"webhook":{"name":"After","url":"https://example.com/update-after","subscriptions":["contact_created"]}}`
w := httptest.NewRecorder()
req, _ := http.NewRequest("PATCH", fmt.Sprintf("/api/v1/accounts/%d/webhooks/%d", s.account.ID, created.ID), bytes.NewBufferString(body))
req.Header.Set("Content-Type", "application/json")
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusOK, w.Code)
var parsed map[string]any
s.Require().NoError(json.Unmarshal(w.Body.Bytes(), &parsed))
webhook := parsed["payload"].(map[string]any)["webhook"].(map[string]any)
s.Equal("After", webhook["name"])
s.Equal("https://example.com/update-after", webhook["url"])
s.Equal([]any{"contact_created"}, webhook["subscriptions"])
}
func (s *WebhookSubscriptionHandlerTestSuite) TestDelete_Success_ReturnsEmptyOK() {
created, err := s.handler.webhookSubscriptionService.CreateWebhook(context.Background(), s.account.ID, service.WebhookSubscriptionMutation{
Name: "Delete",
URL: "https://example.com/delete-hook",
Subscriptions: []string{"message_created"},
})
s.Require().NoError(err)
r := gin.New()
r.DELETE("/api/v1/accounts/:account_id/webhooks/:webhook_id", s.handler.Delete)
w := httptest.NewRecorder()
req, _ := http.NewRequest("DELETE", fmt.Sprintf("/api/v1/accounts/%d/webhooks/%d", s.account.ID, created.ID), nil)
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusOK, w.Code)
s.Empty(w.Body.String())
}
func (s *WebhookSubscriptionHandlerTestSuite) TestCreate_BadRequest_EmptyBody() {
r := gin.New()
r.POST("/api/v1/accounts/:account_id/webhooks/:webhook_id/subscriptions", s.handler.Create)
r.POST("/api/v1/accounts/:account_id/webhooks", s.handler.Create)
w := httptest.NewRecorder()
req, _ := http.NewRequest("POST", fmt.Sprintf("/api/v1/accounts/%d/webhooks/1/subscriptions", s.account.ID), nil)
req, _ := http.NewRequest("POST", fmt.Sprintf("/api/v1/accounts/%d/webhooks", s.account.ID), nil)
req.Header.Set("Content-Type", "application/json")
r.ServeHTTP(w, req)
@@ -78,22 +157,21 @@ func (s *WebhookSubscriptionHandlerTestSuite) TestCreate_BadRequest_EmptyBody()
func (s *WebhookSubscriptionHandlerTestSuite) TestGet_BadRequest_InvalidID() {
r := gin.New()
r.GET("/api/v1/accounts/:account_id/webhooks/:webhook_id/subscriptions/:webhook_id", s.handler.Get)
r.GET("/api/v1/accounts/:account_id/webhooks/:webhook_id", s.handler.Get)
w := httptest.NewRecorder()
req, _ := http.NewRequest("GET", fmt.Sprintf("/api/v1/accounts/%d/webhooks/1/subscriptions/abc", s.account.ID), nil)
req, _ := http.NewRequest("GET", fmt.Sprintf("/api/v1/accounts/%d/webhooks/abc", s.account.ID), nil)
r.ServeHTTP(w, req)
// Get returns 500 for invalid id (parseUintParam error → internal server error path)
assert.NotEqual(s.T(), http.StatusOK, w.Code)
assert.Equal(s.T(), http.StatusBadRequest, w.Code)
}
func (s *WebhookSubscriptionHandlerTestSuite) TestUpdate_BadRequest_InvalidID() {
r := gin.New()
r.PUT("/api/v1/accounts/:account_id/webhooks/:webhook_id/subscriptions/:webhook_id", s.handler.Update)
r.PATCH("/api/v1/accounts/:account_id/webhooks/:webhook_id", s.handler.Update)
w := httptest.NewRecorder()
req, _ := http.NewRequest("PUT", fmt.Sprintf("/api/v1/accounts/%d/webhooks/1/subscriptions/abc", s.account.ID), bytes.NewBufferString(`{"url":"https://example.com"}`))
req, _ := http.NewRequest("PATCH", fmt.Sprintf("/api/v1/accounts/%d/webhooks/abc", s.account.ID), bytes.NewBufferString(`{"webhook":{"url":"https://example.com","subscriptions":["message_created"]}}`))
req.Header.Set("Content-Type", "application/json")
r.ServeHTTP(w, req)
@@ -102,10 +180,10 @@ func (s *WebhookSubscriptionHandlerTestSuite) TestUpdate_BadRequest_InvalidID()
func (s *WebhookSubscriptionHandlerTestSuite) TestDelete_BadRequest_InvalidID() {
r := gin.New()
r.DELETE("/api/v1/accounts/:account_id/webhooks/:webhook_id/subscriptions/:webhook_id", s.handler.Delete)
r.DELETE("/api/v1/accounts/:account_id/webhooks/:webhook_id", s.handler.Delete)
w := httptest.NewRecorder()
req, _ := http.NewRequest("DELETE", fmt.Sprintf("/api/v1/accounts/%d/webhooks/1/subscriptions/abc", s.account.ID), nil)
req, _ := http.NewRequest("DELETE", fmt.Sprintf("/api/v1/accounts/%d/webhooks/abc", s.account.ID), nil)
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusBadRequest, w.Code)
@@ -120,4 +198,4 @@ func (s *WebhookSubscriptionHandlerTestSuite) TestListDeliveries_BadRequest_Inva
r.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusBadRequest, w.Code)
}
}
+28 -24
View File
@@ -10,20 +10,24 @@ import (
// WebhookSubscription represents an account-level webhook subscription for outgoing events.
// Reference: Chatwoot webhook integration + P2B M8 spec
type WebhookSubscription struct {
ID uint `gorm:"primaryKey" json:"id"`
AccountID uint `gorm:"not null;index" json:"account_id"`
URL string `gorm:"size:2048;not null" json:"url"`
Events json.RawMessage `gorm:"type:jsonb;not null" json:"events"` // JSON array of event types, e.g. ["conversation_created","message_created"]
Secret string `gorm:"size:128;not null" json:"secret,omitempty"` // HMAC-SHA256 signing secret
Active bool `gorm:"default:true" json:"active"`
VerifiedAt *time.Time `json:"verified_at,omitempty"`
LastDeliveryStatus string `gorm:"size:50" json:"last_delivery_status,omitempty"` // success/failed/pending
LastDeliveryAt *time.Time `json:"last_delivery_at,omitempty"`
CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"`
UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"`
DeletedAt gorm.DeletedAt `gorm:"index" json:"deleted_at,omitempty"`
ID uint `gorm:"primaryKey" json:"id"`
AccountID uint `gorm:"not null;index" json:"account_id"`
InboxID *uint `gorm:"index" json:"inbox_id,omitempty"`
Name string `gorm:"size:255" json:"name,omitempty"`
URL string `gorm:"size:2048;not null" json:"url"`
Events json.RawMessage `gorm:"type:jsonb;not null" json:"subscriptions"` // JSON array of event types, e.g. ["conversation_created","message_created"]
Secret string `gorm:"size:128;not null" json:"secret,omitempty"` // HMAC-SHA256 signing secret
WebhookType int `gorm:"default:0" json:"webhook_type,omitempty"`
Active bool `gorm:"default:true" json:"active"`
VerifiedAt *time.Time `json:"verified_at,omitempty"`
LastDeliveryStatus string `gorm:"size:50" json:"last_delivery_status,omitempty"` // success/failed/pending
LastDeliveryAt *time.Time `json:"last_delivery_at,omitempty"`
CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"`
UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"`
DeletedAt gorm.DeletedAt `gorm:"index" json:"deleted_at,omitempty"`
Account Account `gorm:"foreignKey:AccountID" json:"account,omitempty"`
Inbox Inbox `gorm:"foreignKey:InboxID" json:"inbox,omitempty"`
}
func (WebhookSubscription) TableName() string { return "webhook_subscriptions" }
@@ -54,19 +58,19 @@ func (s *WebhookSubscription) IsEventSubscribed(eventType string) bool {
// WebhookDelivery represents a single webhook delivery attempt.
// Reference: Chatwoot webhook delivery tracking + P2B M8 spec
type WebhookDelivery struct {
ID uint `gorm:"primaryKey" json:"id"`
SubscriptionID uint `gorm:"not null;index" json:"subscription_id"`
EventType string `gorm:"size:100;not null;index" json:"event_type"`
Payload json.RawMessage `gorm:"type:jsonb" json:"payload"`
ResponseCode int `json:"response_code,omitempty"`
ResponseBody string `gorm:"size:4096" json:"response_body,omitempty"`
Status string `gorm:"size:50;not null;index" json:"status"` // success/failed/retrying
Attempts int `gorm:"default:0" json:"attempts"`
NextAttemptAt *time.Time `json:"next_attempt_at,omitempty"`
CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"`
UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"`
ID uint `gorm:"primaryKey" json:"id"`
SubscriptionID uint `gorm:"not null;index" json:"subscription_id"`
EventType string `gorm:"size:100;not null;index" json:"event_type"`
Payload json.RawMessage `gorm:"type:jsonb" json:"payload"`
ResponseCode int `json:"response_code,omitempty"`
ResponseBody string `gorm:"size:4096" json:"response_body,omitempty"`
Status string `gorm:"size:50;not null;index" json:"status"` // success/failed/retrying
Attempts int `gorm:"default:0" json:"attempts"`
NextAttemptAt *time.Time `json:"next_attempt_at,omitempty"`
CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"`
UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"`
Subscription WebhookSubscription `gorm:"foreignKey:SubscriptionID" json:"subscription,omitempty"`
}
func (WebhookDelivery) TableName() string { return "webhook_deliveries" }
func (WebhookDelivery) TableName() string { return "webhook_deliveries" }
@@ -2,6 +2,7 @@ package repository
import (
"context"
"encoding/json"
"gorm.io/gorm"
@@ -22,7 +23,32 @@ func NewWebhookSubscriptionRepo(db *gorm.DB) *WebhookSubscriptionRepo {
// FindByID retrieves a webhook subscription by primary key.
func (r *WebhookSubscriptionRepo) FindByID(ctx context.Context, id uint) (*model.WebhookSubscription, error) {
var s model.WebhookSubscription
err := r.db.WithContext(ctx).First(&s, id).Error
err := r.db.WithContext(ctx).Preload("Inbox").First(&s, id).Error
if err != nil {
return nil, err
}
return &s, nil
}
// FindByAccountAndID retrieves a webhook subscription scoped to an account.
func (r *WebhookSubscriptionRepo) FindByAccountAndID(ctx context.Context, accountID, id uint) (*model.WebhookSubscription, error) {
var s model.WebhookSubscription
err := r.db.WithContext(ctx).
Preload("Inbox").
Where("account_id = ?", accountID).
First(&s, id).Error
if err != nil {
return nil, err
}
return &s, nil
}
// FindByAccountAndURL retrieves an active webhook by account and URL.
func (r *WebhookSubscriptionRepo) FindByAccountAndURL(ctx context.Context, accountID uint, url string) (*model.WebhookSubscription, error) {
var s model.WebhookSubscription
err := r.db.WithContext(ctx).
Where("account_id = ? AND url = ? AND active = ?", accountID, url, true).
First(&s).Error
if err != nil {
return nil, err
}
@@ -37,7 +63,11 @@ func (r *WebhookSubscriptionRepo) GetByID(ctx context.Context, id uint) (*model.
// ListByAccount retrieves all webhook subscriptions for an account.
func (r *WebhookSubscriptionRepo) ListByAccount(ctx context.Context, accountID uint) ([]model.WebhookSubscription, error) {
var subs []model.WebhookSubscription
err := r.db.WithContext(ctx).Where("account_id = ? AND active = ?", accountID, true).Find(&subs).Error
err := r.db.WithContext(ctx).
Preload("Inbox").
Where("account_id = ? AND active = ?", accountID, true).
Order("id ASC").
Find(&subs).Error
return subs, err
}
@@ -45,20 +75,33 @@ func (r *WebhookSubscriptionRepo) ListByAccount(ctx context.Context, accountID u
func (r *WebhookSubscriptionRepo) ListActiveByAccount(ctx context.Context, accountID uint) ([]model.WebhookSubscription, error) {
var subs []model.WebhookSubscription
err := r.db.WithContext(ctx).
Preload("Inbox").
Where("account_id = ? AND active = ?", accountID, true).
Order("id ASC").
Find(&subs).Error
return subs, err
}
// ListByAccountAndEvent retrieves webhook subscriptions that include a specific event type.
func (r *WebhookSubscriptionRepo) ListByAccountAndEvent(ctx context.Context, accountID uint, eventType string) ([]model.WebhookSubscription, error) {
var subs []model.WebhookSubscription
// Use JSON containment operator for PostgreSQL: events @> '["event_type"]'
err := r.db.WithContext(ctx).
Where("account_id = ? AND active = ?", accountID, true).
Where("events @> ?", eventType).
Find(&subs).Error
return subs, err
subs, err := r.ListActiveByAccount(ctx, accountID)
if err != nil {
return nil, err
}
filtered := make([]model.WebhookSubscription, 0, len(subs))
for _, sub := range subs {
var events []string
if err := json.Unmarshal(sub.Events, &events); err != nil {
continue
}
for _, event := range events {
if event == eventType {
filtered = append(filtered, sub)
break
}
}
}
return filtered, nil
}
// Create inserts a new webhook subscription.
@@ -97,4 +140,4 @@ func (r *WebhookSubscriptionRepo) ListDeliveriesBySubscription(ctx context.Conte
Limit(limit).
Find(&deliveries).Error
return deliveries, err
}
}
+1
View File
@@ -554,6 +554,7 @@ func registerV1Routes(g *gin.RouterGroup, h *Handlers) {
g.POST("/accounts/:account_id/webhooks", h.WebhookSubscription.Create)
g.GET("/accounts/:account_id/webhooks/:webhook_id", h.WebhookSubscription.Get)
g.PUT("/accounts/:account_id/webhooks/:webhook_id", h.WebhookSubscription.Update)
g.PATCH("/accounts/:account_id/webhooks/:webhook_id", h.WebhookSubscription.Update)
g.DELETE("/accounts/:account_id/webhooks/:webhook_id", h.WebhookSubscription.Delete)
// Account routes — scoped with AccountScope middleware (ref: Chatwoot namespace :accounts)
@@ -4,10 +4,13 @@ import (
"context"
"crypto/rand"
"encoding/json"
"errors"
"fmt"
"net/url"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"gorm.io/gorm"
)
// WebhookSubscriptionService provides business logic for managing webhook subscriptions.
@@ -16,6 +19,29 @@ type WebhookSubscriptionService struct {
webhookSubRepo *repository.WebhookSubscriptionRepo
}
var allowedWebhookSubscriptions = map[string]struct{}{
"conversation_status_changed": {},
"conversation_updated": {},
"conversation_created": {},
"contact_created": {},
"contact_updated": {},
"message_created": {},
"message_updated": {},
"webwidget_triggered": {},
"inbox_created": {},
"inbox_updated": {},
"conversation_typing_on": {},
"conversation_typing_off": {},
}
// WebhookSubscriptionMutation is the Chatwoot account webhook create/update payload.
type WebhookSubscriptionMutation struct {
InboxID *uint `json:"inbox_id"`
Name string `json:"name"`
URL string `json:"url"`
Subscriptions []string `json:"subscriptions"`
}
// NewWebhookSubscriptionService creates a new WebhookSubscription service with required dependencies.
func NewWebhookSubscriptionService(webhookSubRepo *repository.WebhookSubscriptionRepo) *WebhookSubscriptionService {
return &WebhookSubscriptionService{
@@ -29,8 +55,23 @@ func (s *WebhookSubscriptionService) ListSubscriptions(ctx context.Context, acco
}
// CreateSubscription creates a new webhook subscription with a generated signing secret.
func (s *WebhookSubscriptionService) CreateSubscription(ctx context.Context, accountID uint, url string, events []string) (*model.WebhookSubscription, error) {
eventsJSON, err := json.Marshal(events)
func (s *WebhookSubscriptionService) CreateSubscription(ctx context.Context, accountID uint, webhookURL string, events []string) (*model.WebhookSubscription, error) {
req := WebhookSubscriptionMutation{URL: webhookURL, Subscriptions: events}
return s.CreateWebhook(ctx, accountID, req)
}
// CreateWebhook creates a Chatwoot-compatible account webhook.
func (s *WebhookSubscriptionService) CreateWebhook(ctx context.Context, accountID uint, req WebhookSubscriptionMutation) (*model.WebhookSubscription, error) {
if err := validateWebhookMutation(req, true); err != nil {
return nil, err
}
if existing, err := s.webhookSubRepo.FindByAccountAndURL(ctx, accountID, req.URL); err == nil && existing.ID != 0 {
return nil, fmt.Errorf("url has already been taken")
} else if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
return nil, err
}
eventsJSON, err := json.Marshal(req.Subscriptions)
if err != nil {
return nil, fmt.Errorf("marshal events: %w", err)
}
@@ -43,7 +84,9 @@ func (s *WebhookSubscriptionService) CreateSubscription(ctx context.Context, acc
sub := &model.WebhookSubscription{
AccountID: accountID,
URL: url,
InboxID: req.InboxID,
Name: req.Name,
URL: req.URL,
Events: eventsJSON,
Secret: secret,
Active: true,
@@ -79,6 +122,51 @@ func (s *WebhookSubscriptionService) UpdateSubscription(ctx context.Context, id
return sub, nil
}
// UpdateWebhook updates a Chatwoot-compatible account webhook scoped to the account.
func (s *WebhookSubscriptionService) UpdateWebhook(ctx context.Context, accountID, id uint, req WebhookSubscriptionMutation) (*model.WebhookSubscription, error) {
if err := validateWebhookMutation(req, false); err != nil {
return nil, err
}
sub, err := s.webhookSubRepo.FindByAccountAndID(ctx, accountID, id)
if err != nil {
return nil, fmt.Errorf("find subscription: %w", err)
}
if req.URL != "" && req.URL != sub.URL {
if existing, err := s.webhookSubRepo.FindByAccountAndURL(ctx, accountID, req.URL); err == nil && existing.ID != sub.ID {
return nil, fmt.Errorf("url has already been taken")
} else if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
return nil, err
}
sub.URL = req.URL
}
sub.Name = req.Name
sub.InboxID = req.InboxID
eventsJSON, err := json.Marshal(req.Subscriptions)
if err != nil {
return nil, fmt.Errorf("marshal events: %w", err)
}
sub.Events = eventsJSON
sub.Active = true
if err := s.webhookSubRepo.Update(ctx, sub); err != nil {
return nil, err
}
return s.webhookSubRepo.FindByAccountAndID(ctx, accountID, id)
}
// DeleteWebhook deletes a Chatwoot-compatible account webhook scoped to the account.
func (s *WebhookSubscriptionService) DeleteWebhook(ctx context.Context, accountID, id uint) error {
if _, err := s.webhookSubRepo.FindByAccountAndID(ctx, accountID, id); err != nil {
return fmt.Errorf("find subscription: %w", err)
}
return s.webhookSubRepo.Delete(ctx, id)
}
// GetWebhook retrieves a Chatwoot-compatible account webhook scoped to the account.
func (s *WebhookSubscriptionService) GetWebhook(ctx context.Context, accountID, id uint) (*model.WebhookSubscription, error) {
return s.webhookSubRepo.FindByAccountAndID(ctx, accountID, id)
}
// DeleteSubscription soft-deletes a webhook subscription.
func (s *WebhookSubscriptionService) DeleteSubscription(ctx context.Context, id uint) error {
return s.webhookSubRepo.Delete(ctx, id)
@@ -101,4 +189,27 @@ func generateWebhookSecret() (string, error) {
return "", err
}
return fmt.Sprintf("%x", b), nil
}
}
func validateWebhookMutation(req WebhookSubscriptionMutation, requireURL bool) error {
if requireURL || req.URL != "" {
u, err := url.ParseRequestURI(req.URL)
if err != nil || u == nil || (u.Scheme != "http" && u.Scheme != "https") || u.Host == "" {
return fmt.Errorf("url is invalid")
}
}
if len(req.Subscriptions) == 0 {
return fmt.Errorf("subscriptions is invalid")
}
seen := map[string]struct{}{}
for _, subscription := range req.Subscriptions {
if _, ok := allowedWebhookSubscriptions[subscription]; !ok {
return fmt.Errorf("subscriptions is invalid")
}
if _, ok := seen[subscription]; ok {
return fmt.Errorf("subscriptions is invalid")
}
seen[subscription] = struct{}{}
}
return nil
}
@@ -0,0 +1,74 @@
package service
import (
"context"
"testing"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/stretchr/testify/require"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"gorm.io/gorm/logger"
)
func TestWebhookSubscriptionServiceChatwootParity(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file::memory:"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
require.NoError(t, err)
require.NoError(t, db.AutoMigrate(&model.Account{}, &model.Inbox{}, &model.WebhookSubscription{}, &model.WebhookDelivery{}))
account := model.Account{Name: "webhook-service-account"}
require.NoError(t, db.Create(&account).Error)
inbox := model.Inbox{AccountID: account.ID, Name: "Support"}
require.NoError(t, db.Create(&inbox).Error)
svc := NewWebhookSubscriptionService(repository.NewWebhookSubscriptionRepo(db))
created, err := svc.CreateWebhook(context.Background(), account.ID, WebhookSubscriptionMutation{
InboxID: &inbox.ID,
Name: "Service hook",
URL: "https://example.com/service-hook",
Subscriptions: []string{"message_created", "conversation_updated"},
})
require.NoError(t, err)
require.NotEmpty(t, created.Secret)
require.Equal(t, "Service hook", created.Name)
require.True(t, created.IsEventSubscribed("message_created"))
updated, err := svc.UpdateWebhook(context.Background(), account.ID, created.ID, WebhookSubscriptionMutation{
Name: "Updated service hook",
URL: "https://example.com/service-hook-updated",
Subscriptions: []string{"contact_created"},
})
require.NoError(t, err)
require.Equal(t, "Updated service hook", updated.Name)
require.True(t, updated.IsEventSubscribed("contact_created"))
require.False(t, updated.IsEventSubscribed("message_created"))
matches, err := repository.NewWebhookSubscriptionRepo(db).ListByAccountAndEvent(context.Background(), account.ID, "contact_created")
require.NoError(t, err)
require.Len(t, matches, 1)
otherAccount := model.Account{Name: "other-webhook-service-account"}
require.NoError(t, db.Create(&otherAccount).Error)
require.Error(t, svc.DeleteWebhook(context.Background(), otherAccount.ID, created.ID))
require.NoError(t, svc.DeleteWebhook(context.Background(), account.ID, created.ID))
}
func TestWebhookSubscriptionServiceValidation(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file::memory:"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
require.NoError(t, err)
require.NoError(t, db.AutoMigrate(&model.WebhookSubscription{}))
svc := NewWebhookSubscriptionService(repository.NewWebhookSubscriptionRepo(db))
_, err = svc.CreateWebhook(context.Background(), 1, WebhookSubscriptionMutation{
URL: "ftp://example.com/hook",
Subscriptions: []string{"message_created"},
})
require.Error(t, err)
_, err = svc.CreateWebhook(context.Background(), 1, WebhookSubscriptionMutation{
URL: "https://example.com/hook",
Subscriptions: []string{"not_allowed"},
})
require.Error(t, err)
}