feat: scope automation and sequence unsubscribes to the caller's organization, require manage_settings to create, edit, enable or delete automations, run each integration action only on its own provider's connection (400 action_provider_mismatch), refuse org-permission gates when no organization service is wired, read /integrations/bookings like contacts, scope lead claims, research runs and CRM list cursors to the organization, and bind contact note edits to the contact in the path

This commit is contained in:
Matthew Meszaros
2026-10-04 02:47:10 -07:00
parent e1c4a91ad1
commit f0ec39febd
31 changed files with 286 additions and 98 deletions
+4 -1
View File
@@ -423,7 +423,8 @@ The `listing` routes publish an app to the [community directory](/api/oauth/#pub
| GET | `/webhooks/:id/deliveries` | `WEBHOOKS` |
| POST | `/webhooks/deliveries/:deliveryId/redeliver` | `WEBHOOKS` |
| GET | `/webhooks/throttle-drops` | `WEBHOOKS` |
| GET/POST/DELETE | `/integrations/*` (except `/integrations/slack/*`, which is [JWT only](#slack)) | `INTEGRATIONS` |
| GET/POST/DELETE | `/integrations/*` (except `/integrations/slack/*`, which is [JWT only](#slack), and `/integrations/bookings`) | `INTEGRATIONS` |
| GET | `/integrations/bookings` | `READ_CONTACTS` |
| GET/PUT | `/integrations/salesforce/:id/settings` | `INTEGRATIONS` |
| GET | `/integrations/salesforce/:id/overview`, `/metadata`, `/users`, `/list-views`, `/campaigns`, `/activity` | `INTEGRATIONS` |
| POST | `/integrations/salesforce/:id/import/preview` | `INTEGRATIONS` |
@@ -435,6 +436,8 @@ The `listing` routes publish an app to the [community directory](/api/oauth/#pub
| DELETE | `/contacts/:id/salesforce/links/:linkId` | `INTEGRATIONS` |
| GET/POST/PATCH/DELETE | `/automations[/:id]` | `INTEGRATIONS` |
| PATCH | `/automations/:id/layout` | `INTEGRATIONS` |
A signed-in member needs `manage_settings` to create, edit, enable or delete an automation, because a flow runs as the workspace; `use_integrations` is enough to list, open, test and read the run history.
| GET/POST/PATCH/DELETE | `/warmup/routing[/:id]` | `WARMUP_ROUTING` |
### Retry safety
+6
View File
@@ -331,6 +331,12 @@ The [Salesforce](/guides/salesforce/) endpoints answer with their own codes. Whe
| `import_running` | 409 | `POST /integrations/salesforce/:id/import-sources/:sourceId/run` while that import is already running |
| `salesforce_rate_limited` | 429 | The Salesforce org has used its API requests for today |
#### Automation refusals
| `code` | Status | Meaning |
|--------|--------|---------|
| `action_provider_mismatch` | 400 | `POST /automations`, `PATCH /automations/:id` or `POST /integrations/connections/:id/events` paired an action with a connection whose provider does not run it, such as a Slack message on a HubSpot connection. Each action runs only on the integration it belongs to; the provider's actions are listed under `capability.actions` in `GET /integrations/catalog` |
### 401 Unauthorized
Returned when authentication fails.
@@ -628,7 +628,7 @@ Creating an import is not deduplicated by request: a retried create saves a seco
Returns up to 50 recent booked meetings. For the full Meetings page list with filters and pagination, use `GET /meetings`.
Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_settings` or `use_integrations`).
Auth: **Scope** `READ_CONTACTS` · **Org permission** `view_contacts`.
### Response
@@ -685,7 +685,7 @@ Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_setti
Creates a new automation flow: a trigger event plus a graph of condition and action nodes.
Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_settings` or `use_integrations`).
Auth: **Scope** `INTEGRATIONS` · **Org permission** `manage_settings`.
### Request body
@@ -776,7 +776,7 @@ Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_setti
Replaces an automation's name, enabled state, trigger, filter, and graph. The body shape matches the create payload.
Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_settings` or `use_integrations`).
Auth: **Scope** `INTEGRATIONS` · **Org permission** `manage_settings`.
| Parameter | In | Type | Description |
|-----------|-----|------|-------------|
@@ -816,7 +816,7 @@ Same fields as [create an automation](#create-an-automation).
Removes an automation. Returns `409 Conflict` when the automation is still referenced by campaign steps.
Auth: **Scope** `INTEGRATIONS` (or session with **Org permission** `manage_settings` or `use_integrations`).
Auth: **Scope** `INTEGRATIONS` · **Org permission** `manage_settings`.
| Parameter | In | Type | Description |
|-----------|-----|------|-------------|
+2 -2
View File
@@ -58,14 +58,14 @@ Roles are workspace data. Every workspace starts with three seeded roles that ar
| | View / manage contacts | Read contacts, segments, labels / create, edit, delete |
| | Manage sequences | Edit step content and spacing |
| | View analytics | Deliverability and engagement reports |
| | Use integrations | Push contacts and deals to connected tools |
| | Use integrations | Push contacts and deals to connected tools, view and test automations |
| | Use AI | Remie and AI drafting, spending shared credits |
| **People** | Manage team | Invite, remove, re-role members |
| | Transfer ownership | Hand ownership to another member |
| **Sending** | Manage mailboxes | Connect, disconnect, configure senders |
| | Send campaigns | Start, pause, resume |
| | Use unified inbox | Read and reply from the shared inbox |
| **Workspace** | Manage settings | Workspace-wide settings |
| **Workspace** | Manage settings | Workspace-wide settings, connecting integrations, building automations |
| | Manage billing | Checkout, invoices, plan changes and cancellation |
| | Manage API keys | Create and revoke workspace keys |
+1 -1
View File
@@ -19567,7 +19567,7 @@
"get": {
"operationId": "integrations_bookings_list",
"summary": "List meeting bookings (integrations view)",
"description": "Up to 50 recent booked meetings, surfaced on the integrations page. For the full Meetings list with filters and pagination, use `GET /meetings`.",
"description": "Up to 50 recent booked meetings, surfaced on the integrations page. For the full Meetings list with filters and pagination, use `GET /meetings`. Scope `READ_CONTACTS`.",
"tags": [
"integrations"
],
+12 -2
View File
@@ -94,6 +94,11 @@ func (h *Handler) UpdateContactNote(c *gin.Context) {
errx.Handle(c, errx.New(errx.BadRequest, "no organization selected"))
return
}
contactID, err := uuid.Parse(c.Param("id"))
if err != nil {
errx.Handle(c, errx.ErrUuid)
return
}
noteID, err := uuid.Parse(c.Param("noteId"))
if err != nil {
errx.Handle(c, errx.ErrUuid)
@@ -106,7 +111,7 @@ func (h *Handler) UpdateContactNote(c *gin.Context) {
return
}
note, xerr := h.CRMService.UpdateNote(c.Request.Context(), *orgID, noteID, &data)
note, xerr := h.CRMService.UpdateNote(c.Request.Context(), *orgID, &contactID, noteID, &data)
if xerr != nil {
errx.Handle(c, xerr)
return
@@ -123,13 +128,18 @@ func (h *Handler) DeleteContactNote(c *gin.Context) {
errx.Handle(c, errx.New(errx.BadRequest, "no organization selected"))
return
}
contactID, err := uuid.Parse(c.Param("id"))
if err != nil {
errx.Handle(c, errx.ErrUuid)
return
}
noteID, err := uuid.Parse(c.Param("noteId"))
if err != nil {
errx.Handle(c, errx.ErrUuid)
return
}
xerr := h.CRMService.DeleteNote(c.Request.Context(), *orgID, noteID)
xerr := h.CRMService.DeleteNote(c.Request.Context(), *orgID, &contactID, noteID)
if xerr != nil {
errx.Handle(c, xerr)
return
+11 -3
View File
@@ -337,6 +337,14 @@ func (h *Handler) ListConnectionEventSubscriptions(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"events": subs})
}
// automationWriteError maps a refused automation or subscription write to a 400.
func automationWriteError(err error) *errx.Error {
if errors.Is(err, integration.ErrActionProviderMismatch) {
return errx.NewWithIdentifier(errx.BadRequest, "action_provider_mismatch", err.Error())
}
return errx.New(errx.BadRequest, err.Error())
}
func (h *Handler) CreateConnectionEventSubscription(c *gin.Context) {
orgID, userID, ok := h.requireIntegrationActor(c, true)
if !ok {
@@ -359,7 +367,7 @@ func (h *Handler) CreateConnectionEventSubscription(c *gin.Context) {
sub, err := h.IntegrationService.CreateEventSubscription(c.Request.Context(), orgID, connID,
strings.TrimSpace(p.EventType), models.IntegrationAction(strings.TrimSpace(p.Action)), p.Config, enabled)
if err != nil {
errx.JSON(c, errx.New(errx.BadRequest, err.Error()))
errx.JSON(c, automationWriteError(err))
return
}
h.auditIntegration(c, userID, models.AuditActionUpdate, connID, "event:"+p.EventType)
@@ -812,7 +820,7 @@ func (h *Handler) CreateAutomation(c *gin.Context) {
}
a, err := h.IntegrationService.CreateAutomation(c.Request.Context(), orgID, w)
if err != nil {
errx.JSON(c, errx.New(errx.BadRequest, err.Error()))
errx.JSON(c, automationWriteError(err))
return
}
h.auditIntegrationEntity(c, userID, models.AuditActionCreate, models.AuditEntityAutomation, a.ID, a.Name)
@@ -837,7 +845,7 @@ func (h *Handler) UpdateAutomation(c *gin.Context) {
}
a, err := h.IntegrationService.UpdateAutomation(c.Request.Context(), orgID, id, w)
if err != nil {
errx.JSON(c, errx.New(errx.BadRequest, err.Error()))
errx.JSON(c, automationWriteError(err))
return
}
h.auditIntegrationEntity(c, userID, models.AuditActionUpdate, models.AuditEntityAutomation, id, a.Name)
+5 -3
View File
@@ -238,9 +238,10 @@ func (h *Handler) RequireAccess(orgPerm models.OrganizationPermission, apiPerm u
}
c.Next()
default:
// JWT path: defer to the org-permission gate.
// JWT path: defer to the org-permission gate, which refuses when it cannot check.
if h.OrganizationService == nil {
c.Next()
errx.JSON(c, errx.InternalError())
c.Abort()
return
}
userID, err := GetUserUUID(c)
@@ -308,7 +309,8 @@ func (h *Handler) RequireAnyAccess(apiPerm uint64, orgPerms ...models.Organizati
c.Next()
default:
if h.OrganizationService == nil {
c.Next()
errx.JSON(c, errx.InternalError())
c.Abort()
return
}
userID, err := GetUserUUID(c)
+26
View File
@@ -178,3 +178,29 @@ func TestRequireAccessWithQueryOnlyGatesWithTheParam(t *testing.T) {
})
}
}
// A session caller is refused when no organization service is wired to check it.
func TestOrgGatesRefuseWithoutOrganizationService(t *testing.T) {
h := &Handler{}
for name, gate := range map[string]gin.HandlerFunc{
"RequireAccess": h.RequireAccess(models.PermManageSettings, models.APIPermIntegrations),
"RequireAnyAccess": h.RequireAnyAccess(models.APIPermIntegrations, models.PermManageSettings, models.PermUseIntegrations),
"RequirePermission": h.RequirePermission(models.PermManageSettings),
} {
r := gin.New()
r.Use(func(c *gin.Context) {
c.Set(AuthTypeKey, AuthTypeJWT)
c.Next()
})
calls := 0
r.GET("/x", gate, func(c *gin.Context) {
calls++
c.JSON(http.StatusOK, gin.H{})
})
w := httptest.NewRecorder()
r.ServeHTTP(w, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/x", nil))
if w.Code != http.StatusInternalServerError || calls != 0 {
t.Errorf("%s: status = %d, handler ran %d times; want 500 and 0", name, w.Code, calls)
}
}
}
+2 -1
View File
@@ -174,7 +174,8 @@ func (h *Handler) RequireOrganization() gin.HandlerFunc {
func (h *Handler) RequirePermission(perm models.OrganizationPermission) gin.HandlerFunc {
return func(c *gin.Context) {
if h.OrganizationService == nil {
c.Next()
errx.JSON(c, errx.InternalError())
c.Abort()
return
}
+5 -6
View File
@@ -1174,7 +1174,8 @@ func Run(
integrations.POST("/connections/:id/rotate-inbound-url", write, h.RotateConnectionInboundURL)
integrations.POST("/connections/:id/test", write, h.TestConnection)
integrations.POST("/connections/:id/push", operate, h.PushContactsToIntegration)
integrations.GET("/bookings", read, h.ListMeetingBookings)
// Bookings carry invitee details, so they are read like contacts.
integrations.GET("/bookings", m.RequireAccess(models.PermViewContacts, models.APIPermReadContacts), h.ListMeetingBookings)
// Native Salesforce sync (:id is the connection). Reading health and
// the activity log is operational; changing what syncs is settings.
@@ -1213,15 +1214,13 @@ func Run(
// Automations (org-scoped). The visual flow builder: a trigger event +
// action steps across integrations. Reads reachable by operational
// integration users; creating/editing is a settings action.
// integration users; creating/editing is a settings action, because a
// flow runs as the workspace.
automations := protected.Group("/automations")
automations.Use(m.RequireOrganization(), m.RateLimitMiddleware(models.RateLimitWrite))
{
aread := m.RequireAnyAccess(models.APIPermIntegrations, models.PermManageSettings, models.PermUseIntegrations)
// Writing automations needs the integration permission (same family as
// reads) OR settings-manager; previously it required manage-settings only,
// which let integration-permitted members open the builder but 403 on save.
awrite := m.RequireAnyAccess(models.APIPermIntegrations, models.PermManageSettings, models.PermUseIntegrations)
awrite := m.RequireAccess(models.PermManageSettings, models.APIPermIntegrations)
automations.GET("", aread, h.ListAutomations)
automations.POST("", awrite, h.CreateAutomation)
automations.GET("/:id", aread, h.GetAutomation)
@@ -71,6 +71,13 @@ func (r incomingReplyContactRepo) GetByID(context.Context, uuid.UUID) (*models.C
return r.taskContact, nil
}
func (r incomingReplyContactRepo) GetByIDsAndOrganization(_ context.Context, _ uuid.UUID, ids []uuid.UUID) ([]models.Contact, *errx.Error) {
if r.taskContact == nil || len(ids) != 1 || ids[0] != r.taskContact.ID {
return nil, nil
}
return []models.Contact{*r.taskContact}, nil
}
type incomingReplyProgressRepo struct {
repository.CampaignProgressRepository
replied int
@@ -125,8 +125,8 @@ func TestUnsubscribeLinkOptsOutTheLeadsCopies(t *testing.T) {
{"link", func(s *service, org, campaign, lead uuid.UUID) *errx.Error {
return s.UnsubscribeFromLink(context.Background(), org, campaign, lead, "one_click")
}, []string{"task-contact@example.test", "jonas@acme.test"}},
{"sequence action", func(s *service, _, campaign, lead uuid.UUID) *errx.Error {
return s.Unsubscribe(context.Background(), campaign, lead)
{"sequence action", func(s *service, org, campaign, lead uuid.UUID) *errx.Error {
return s.Unsubscribe(context.Background(), org, campaign, lead)
}, []string{"task-contact@example.test"}},
} {
svc, _, adv, _ := newCopyReplyService(t)
+4 -1
View File
@@ -290,7 +290,10 @@ func (s *service) executeInstantActionNode(ctx context.Context, campaign *models
s.logActionErr(campaign, contact, cfg.Type, eventKind, xerr)
}
case "unsubscribe":
if xerr := s.Unsubscribe(ctx, campaign.ID, contact.ID); xerr != nil {
if campaign.OrganizationID == nil {
return
}
if xerr := s.Unsubscribe(ctx, *campaign.OrganizationID, campaign.ID, contact.ID); xerr != nil {
s.logActionErr(campaign, contact, cfg.Type, eventKind, xerr)
}
case "create_task":
+11 -12
View File
@@ -67,7 +67,7 @@ type Service interface {
// Unsubscribe suppresses a contact in response to a List-Unsubscribe action
// (one-click POST or the manual link). Always suppresses — it's an explicit
// recipient request, independent of the auto-suppress settings.
Unsubscribe(ctx context.Context, campaignID, contactID uuid.UUID) *errx.Error
Unsubscribe(ctx context.Context, organizationID, campaignID, contactID uuid.UUID) *errx.Error
// UnsubscribeFromLink is Unsubscribe for a verified link token: the
// organization in the token must own the campaign, and via names the
// mechanism ("one_click" for the RFC 8058 POST, "link" for a click).
@@ -594,32 +594,31 @@ func (s *service) ListPipelines(ctx context.Context, orgID uuid.UUID) ([]models.
return s.crmRepo.ListPipelines(ctx, orgID)
}
func (s *service) Unsubscribe(ctx context.Context, campaignID, contactID uuid.UUID) *errx.Error {
return s.unsubscribe(ctx, nil, campaignID, contactID, "action")
func (s *service) Unsubscribe(ctx context.Context, organizationID, campaignID, contactID uuid.UUID) *errx.Error {
return s.unsubscribe(ctx, organizationID, campaignID, contactID, "action")
}
func (s *service) UnsubscribeFromLink(ctx context.Context, organizationID, campaignID, contactID uuid.UUID, via string) *errx.Error {
if via != "one_click" {
via = "link"
}
return s.unsubscribe(ctx, &organizationID, campaignID, contactID, via)
return s.unsubscribe(ctx, organizationID, campaignID, contactID, via)
}
// unsubscribe records an explicit opt-out: the address goes on the workspace
// suppression list and the contact's own subscription flag is cleared, so the
// CRM and the send gate tell the same story.
func (s *service) unsubscribe(ctx context.Context, expectOrg *uuid.UUID, campaignID, contactID uuid.UUID, via string) *errx.Error {
// CRM and the send gate tell the same story. The campaign and the contact
// must both belong to organizationID.
func (s *service) unsubscribe(ctx context.Context, organizationID, campaignID, contactID uuid.UUID, via string) *errx.Error {
campaign, err := s.campaignRepo.GetByID(ctx, campaignID)
if err != nil || campaign == nil || campaign.OrganizationID == nil {
if err != nil || campaign == nil || campaign.OrganizationID == nil || *campaign.OrganizationID != organizationID {
return errx.New(errx.BadRequest, "invalid unsubscribe link")
}
if expectOrg != nil && *expectOrg != *campaign.OrganizationID {
return errx.New(errx.BadRequest, "invalid unsubscribe link")
}
contact, cerr := s.contactRepo.GetByID(ctx, contactID)
if cerr != nil || contact == nil || contact.Email == "" {
found, cerr := s.contactRepo.GetByIDsAndOrganization(ctx, organizationID, []uuid.UUID{contactID})
if cerr != nil || len(found) != 1 || found[0].Email == "" {
return errx.New(errx.BadRequest, "invalid unsubscribe link")
}
contact := &found[0]
reason := map[string]string{
"one_click": "one-click unsubscribe (mail client)",
+2 -2
View File
@@ -603,7 +603,7 @@ func (d Deps) updateContactNote(ctx context.Context, inv Invocation, args json.R
if in.Body == "" {
return "", ErrInvalidArgs
}
note, xerr := d.CRM.UpdateNote(ctx, inv.OrgID, nid, &models.UpdateContactNote{Content: &in.Body})
note, xerr := d.CRM.UpdateNote(ctx, inv.OrgID, nil, nid, &models.UpdateContactNote{Content: &in.Body})
if xerr != nil {
return "", fromErrx(xerr)
}
@@ -622,7 +622,7 @@ func (d Deps) deleteContactNote(ctx context.Context, inv Invocation, args json.R
if err != nil {
return "", err
}
if xerr := d.CRM.DeleteNote(ctx, inv.OrgID, nid); xerr != nil {
if xerr := d.CRM.DeleteNote(ctx, inv.OrgID, nil, nid); xerr != nil {
return "", fromErrx(xerr)
}
d.logAudit(ctx, inv, models.AuditActionDelete, models.AuditEntityCRMNote, &nid, nil)
+26 -4
View File
@@ -16,8 +16,9 @@ type CRMService interface {
// Notes
CreateNote(ctx context.Context, orgID, contactID, userID uuid.UUID, data *models.CreateContactNote) (*models.ContactNote, *errx.Error)
ListNotes(ctx context.Context, orgID, contactID uuid.UUID, limit int, cursor *uuid.UUID) (*models.ContactNotesResult, *errx.Error)
UpdateNote(ctx context.Context, orgID, noteID uuid.UUID, data *models.UpdateContactNote) (*models.ContactNote, *errx.Error)
DeleteNote(ctx context.Context, orgID, noteID uuid.UUID) *errx.Error
// UpdateNote and DeleteNote refuse a note on another contact when contactID is set.
UpdateNote(ctx context.Context, orgID uuid.UUID, contactID *uuid.UUID, noteID uuid.UUID, data *models.UpdateContactNote) (*models.ContactNote, *errx.Error)
DeleteNote(ctx context.Context, orgID uuid.UUID, contactID *uuid.UUID, noteID uuid.UUID) *errx.Error
// Activities
ListActivities(ctx context.Context, orgID, contactID uuid.UUID, limit int, cursor *uuid.UUID) (*models.ContactActivitiesResult, *errx.Error)
@@ -170,13 +171,31 @@ func (s *crmService) ListNotes(ctx context.Context, orgID, contactID uuid.UUID,
return result, nil
}
func (s *crmService) UpdateNote(ctx context.Context, orgID, noteID uuid.UUID, data *models.UpdateContactNote) (*models.ContactNote, *errx.Error) {
// noteOnContact checks that noteID is a note on contactID, when one is given.
func (s *crmService) noteOnContact(ctx context.Context, orgID uuid.UUID, contactID *uuid.UUID, noteID uuid.UUID) *errx.Error {
if contactID == nil {
return nil
}
note, err := s.repo.GetNote(ctx, orgID, noteID)
if err != nil {
return toErrx(err)
}
if note.ContactID != *contactID {
return errx.ErrNotFound
}
return nil
}
func (s *crmService) UpdateNote(ctx context.Context, orgID uuid.UUID, contactID *uuid.UUID, noteID uuid.UUID, data *models.UpdateContactNote) (*models.ContactNote, *errx.Error) {
if data.Content == nil || len(*data.Content) == 0 {
return nil, errx.New(errx.BadRequest, "content is required")
}
if len(*data.Content) > 10000 {
return nil, errx.New(errx.BadRequest, "content must be at most 10000 characters")
}
if xerr := s.noteOnContact(ctx, orgID, contactID, noteID); xerr != nil {
return nil, xerr
}
if ext := s.external(ctx, orgID); ext != nil {
if xerr := ext.PushNoteUpdate(ctx, orgID, noteID, *data.Content); xerr != nil {
@@ -196,7 +215,10 @@ func (s *crmService) UpdateNote(ctx context.Context, orgID, noteID uuid.UUID, da
return note, nil
}
func (s *crmService) DeleteNote(ctx context.Context, orgID, noteID uuid.UUID) *errx.Error {
func (s *crmService) DeleteNote(ctx context.Context, orgID uuid.UUID, contactID *uuid.UUID, noteID uuid.UUID) *errx.Error {
if xerr := s.noteOnContact(ctx, orgID, contactID, noteID); xerr != nil {
return xerr
}
if ext := s.external(ctx, orgID); ext != nil {
if xerr := ext.PushNoteDelete(ctx, orgID, noteID); xerr != nil {
return xerr
@@ -0,0 +1,42 @@
package integration
import (
"context"
"errors"
"testing"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
func TestProviderSupportsAction(t *testing.T) {
for _, tc := range []struct {
provider models.IntegrationProvider
action models.IntegrationAction
want bool
}{
{models.IntegrationSlack, models.IntegrationActionSlackNotify, true},
{models.IntegrationHubSpot, models.IntegrationActionHubSpotUpsert, true},
{models.IntegrationZapier, models.IntegrationActionGenericWebhookPing, true},
{models.IntegrationHubSpot, models.IntegrationActionSlackNotify, false},
{models.IntegrationSlack, models.IntegrationActionGenericWebhookPing, false},
{models.IntegrationCalendly, models.IntegrationActionDiscordNotify, false},
{models.IntegrationProvider("unknown"), models.IntegrationActionSlackNotify, false},
} {
if got := models.ProviderSupportsAction(tc.provider, tc.action); got != tc.want {
t.Errorf("ProviderSupportsAction(%s, %s) = %v, want %v", tc.provider, tc.action, got, tc.want)
}
}
}
// execAction refuses before it opens a connection's secrets for another provider's action.
func TestExecActionRefusesAnotherProvidersAction(t *testing.T) {
target := repository.DispatchTarget{
Subscription: models.IntegrationEventSubscription{Action: models.IntegrationActionSlackNotify},
Secrets: repository.ConnectionSecrets{Conn: models.IntegrationConnection{Provider: models.IntegrationHubSpot}},
}
err := (&service{}).execAction(context.Background(), target, map[string]any{})
if !errors.Is(err, ErrActionProviderMismatch) {
t.Fatalf("execAction = %v, want ErrActionProviderMismatch", err)
}
}
+7
View File
@@ -21,6 +21,9 @@ import (
// the dispatcher must not overwrite that status when it sees this error.
var errReauthRequired = errors.New("reauth required")
// ErrActionProviderMismatch means the action does not run on the connection's provider.
var ErrActionProviderMismatch = errors.New("this action does not run on the selected integration")
// Dispatch fans a platform event out to every matching event subscription.
// Targets are resolved synchronously (cheap, indexed) but the provider calls
// run on a detached context so the caller (an API handler or consumer) never
@@ -110,6 +113,10 @@ func (s *service) runAction(ctx context.Context, target repository.DispatchTarge
// upserts so behaviour follows the user's configuration instead of a fixed shape.
func (s *service) execAction(ctx context.Context, target repository.DispatchTarget, data map[string]any) error {
sub := target.Subscription
// A connection's credentials only ever reach its own provider.
if !models.ProviderSupportsAction(target.Secrets.Conn.Provider, sub.Action) {
return ErrActionProviderMismatch
}
secretCfg, err := s.openConfig(ctx, &target.Secrets)
if err != nil {
return fmt.Errorf("decrypt config: %w", err)
+2 -2
View File
@@ -169,7 +169,7 @@ type NativeActions interface {
CreateTask(ctx context.Context, orgID, createdBy uuid.UUID, data *models.CreateCRMTask) error
CreateDeal(ctx context.Context, orgID, createdBy uuid.UUID, data *models.CreateDeal) error
MoveDealStage(ctx context.Context, orgID, contactID, pipelineID, stageID uuid.UUID) error
Unsubscribe(ctx context.Context, campaignID, contactID uuid.UUID) error
Unsubscribe(ctx context.Context, orgID, campaignID, contactID uuid.UUID) error
// LabelThread additively applies unibox conversation labels to a thread, on
// behalf of the mailbox-owner userID (categories are per user). Backs the
// "label_email" action; userID + threadID come from the reply event data.
@@ -504,7 +504,7 @@ func (s *service) execNativeAction(ctx context.Context, a models.Automation, n m
if perr != nil {
return fmt.Errorf("unsubscribe needs a campaign_id in the event data")
}
return s.native.Unsubscribe(ctx, campID, c.ID)
return s.native.Unsubscribe(ctx, a.OrganizationID, campID, c.ID)
case models.IntegrationActionAddTag, models.IntegrationActionRemoveTag:
catID, perr := uuid.Parse(cfg.CategoryID)
+6
View File
@@ -687,6 +687,9 @@ func (s *service) CreateEventSubscription(ctx context.Context, orgID, connID uui
if !models.IsValidWebhookEventType(eventType) {
return nil, fmt.Errorf("unknown event type: %s", eventType)
}
if !models.ProviderSupportsAction(conn.Provider, action) {
return nil, ErrActionProviderMismatch
}
// SSRF guard for action configs that carry an outbound URL.
if err := validateOutboundConfigURLs(config); err != nil {
return nil, err
@@ -911,6 +914,9 @@ func (s *service) validateAutomationGraph(ctx context.Context, orgID uuid.UUID,
if conn == nil {
return errors.New("an action node references an unknown integration")
}
if !models.ProviderSupportsAction(conn.Provider, n.Action) {
return ErrActionProviderMismatch
}
cfg := map[string]any{}
if len(n.Config) > 0 {
_ = json.Unmarshal(n.Config, &cfg)
+2 -2
View File
@@ -110,8 +110,8 @@ func (a Adapter) MoveDealStage(ctx context.Context, orgID, contactID, pipelineID
return nil
}
func (a Adapter) Unsubscribe(ctx context.Context, campaignID, contactID uuid.UUID) error {
if e := a.Adv.Unsubscribe(ctx, campaignID, contactID); e != nil {
func (a Adapter) Unsubscribe(ctx context.Context, orgID, campaignID, contactID uuid.UUID) error {
if e := a.Adv.Unsubscribe(ctx, orgID, campaignID, contactID); e != nil {
return e
}
return nil
+3
View File
@@ -113,6 +113,9 @@ func (s *service) RunResearch(ctx context.Context, inv aitools.Invocation, conta
Status: models.ResearchRunning, Objective: strings.TrimSpace(objective),
}
created, replayed, err := s.repo.CreateRun(ctx, run, idempotencyKey)
if errors.Is(err, repository.ErrResearchContactNotFound) {
return nil, errx.New(errx.NotFound, "contact not found")
}
if err != nil {
return nil, errx.New(errx.Internal, "failed to create research run")
}
@@ -94,6 +94,12 @@ func warmblyContactFields() []FieldDef {
}
}
// ProviderSupportsAction reports whether action runs on a connection of provider.
func ProviderSupportsAction(provider IntegrationProvider, action IntegrationAction) bool {
c := CapabilityFor(provider)
return c != nil && c.Action(action) != nil
}
// CapabilityFor returns the descriptor for a provider, or nil if the provider
// has no configurable capability surface.
func CapabilityFor(provider IntegrationProvider) *ProviderCapability {
+4 -3
View File
@@ -452,13 +452,14 @@ func (r *contactRepository) Add(ctx context.Context, userID string, orgID uuid.U
}
campaignLinks = append(campaignLinks, added...)
// A hand-picked add ends a manual removal, as the bulk add does.
if _, err := tx.Exec(ctx, `DELETE FROM campaign_lead_removals WHERE contact_id = $1 AND campaign_id = ANY($2)`, ncontacts[i].ID, cids); err != nil {
if _, err := tx.Exec(ctx, `DELETE FROM campaign_lead_removals WHERE contact_id = $1 AND campaign_id = ANY($2)
AND campaign_id IN (SELECT id FROM campaigns WHERE organization_id = $3)`, ncontacts[i].ID, cids, orgID); err != nil {
db.CaptureError(err, "", nil, "campaign_lead_removals clear")
return nil, errx.InternalError()
}
// It also claims a lead a linked segment had enrolled, so detaching
// that segment later does not withdraw somebody's hand-picked lead.
if _, err := tx.Exec(ctx, claimLeadsManualSQL, cids, []uuid.UUID{ncontacts[i].ID}); err != nil {
if _, err := tx.Exec(ctx, claimLeadsManualSQL, cids, []uuid.UUID{ncontacts[i].ID}, orgID); err != nil {
db.CaptureError(err, "", nil, "campaign_leads claim")
return nil, errx.InternalError()
}
@@ -2912,7 +2913,7 @@ func (r *contactRepository) BulkUpdate(ctx context.Context, userID string, orgID
}
// Claims leads a linked segment had enrolled, so detaching that
// segment later leaves hand-picked leads alone.
if _, err := tx.Exec(ctx, claimLeadsManualSQL, data.AddCampaigns, data.Contacts); err != nil {
if _, err := tx.Exec(ctx, claimLeadsManualSQL, data.AddCampaigns, data.Contacts, orgID); err != nil {
db.CaptureError(err, "", nil, "campaign_leads claim")
return nil, errx.InternalError()
}
+4 -4
View File
@@ -129,7 +129,7 @@ func (r *crmRepository) ListNotes(ctx context.Context, orgID, contactID uuid.UUI
WHERE cn.contact_id = $1
AND cn.organization_id = $4
AND ($2::uuid IS NULL OR (cn.created_at, cn.id) < (
SELECT created_at, id FROM contact_notes WHERE id = $2
SELECT created_at, id FROM contact_notes WHERE id = $2 AND organization_id = $4
))
ORDER BY cn.created_at DESC, cn.id DESC
LIMIT $3
@@ -222,7 +222,7 @@ func (r *crmRepository) ListActivities(ctx context.Context, orgID, contactID uui
WHERE contact_id = $1
AND organization_id = $4
AND ($2::uuid IS NULL OR (created_at, id) < (
SELECT created_at, id FROM contact_activities WHERE id = $2
SELECT created_at, id FROM contact_activities WHERE id = $2 AND organization_id = $4
))
ORDER BY created_at DESC, id DESC
LIMIT $3
@@ -699,7 +699,7 @@ func (r *crmRepository) ListDeals(ctx context.Context, orgID uuid.UUID, pipeline
}
if cursor != nil {
whereClauses = append(whereClauses, fmt.Sprintf("(created_at, id) < (SELECT created_at, id FROM deals WHERE id = $%d)", argPos))
whereClauses = append(whereClauses, fmt.Sprintf("(created_at, id) < (SELECT created_at, id FROM deals WHERE id = $%d AND organization_id = $1)", argPos))
args = append(args, *cursor)
argPos++
}
@@ -1233,7 +1233,7 @@ func (r *crmRepository) ListCRMTasks(ctx context.Context, orgID uuid.UUID, conta
}
if cursor != nil {
whereClauses = append(whereClauses, fmt.Sprintf("(created_at, id) < (SELECT created_at, id FROM crm_tasks WHERE id = $%d)", argPos))
whereClauses = append(whereClauses, fmt.Sprintf("(created_at, id) < (SELECT created_at, id FROM crm_tasks WHERE id = $%d AND organization_id = $1)", argPos))
args = append(args, *cursor)
argPos++
}
+10 -1
View File
@@ -11,6 +11,9 @@ import (
"github.com/warmbly/warmbly/internal/models"
)
// ErrResearchContactNotFound means the run's contact is not in its organization.
var ErrResearchContactNotFound = errors.New("contact not found")
// ResearchRepository persists contact research runs.
type ResearchRepository interface {
// CreateRun inserts a run. If idempotencyKey is non-empty and a run already
@@ -69,12 +72,18 @@ func (r *researchRepository) CreateRun(ctx context.Context, run *models.ContactR
if idempotencyKey != "" {
keyArg = &idempotencyKey
}
// The contact must belong to the run's organization; no row means it does not.
out := &models.ContactResearchRun{}
err = scanRun(r.DB.QueryRow(ctx, `
INSERT INTO contact_research_runs (org_id, contact_id, requested_by, status, objective, result, idempotency_key)
VALUES ($1, $2, $3, $4, $5, $6, $7)
SELECT $1, c.id, $3, $4, $5, $6, $7
FROM contacts c
WHERE c.id = $2 AND c.organization_id = $1
RETURNING `+researchCols,
run.OrgID, run.ContactID, run.RequestedBy, run.Status, run.Objective, resultRaw, keyArg), out)
if errors.Is(err, pgx.ErrNoRows) {
return nil, false, ErrResearchContactNotFound
}
if err != nil {
return nil, false, err
}
+5 -2
View File
@@ -433,9 +433,12 @@ const (
// claimLeadsManualSQL promotes leads a person chose to 'manual'. The insert
// paths cannot do this in their own ON CONFLICT clause: DO UPDATE would put
// rows that were already leads into RETURNING, and RETURNING is what writes
// the "added to campaign" activity. $1 is the campaign ids, $2 the contacts.
// the "added to campaign" activity. $1 is the campaign ids, $2 the contacts,
// $3 the organization both must belong to.
const claimLeadsManualSQL = `UPDATE campaign_leads SET source = 'manual'
WHERE campaign_id = ANY($1::uuid[]) AND contact_id = ANY($2::uuid[]) AND source <> 'manual'`
WHERE campaign_id = ANY($1::uuid[]) AND contact_id = ANY($2::uuid[]) AND source <> 'manual'
AND campaign_id IN (SELECT id FROM campaigns WHERE organization_id = $3)
AND contact_id IN (SELECT id FROM contacts WHERE organization_id = $3)`
// insertSegmentLeads enrols every contact matching the precompiled segment
// clause as a lead, logging a campaign_added activity for each row that was
+4 -1
View File
@@ -1428,7 +1428,10 @@ func (s *tasksService) executeActionNode(ctx context.Context, campaign *models.C
}
return nil
case "unsubscribe":
if xerr := s.advanced.Unsubscribe(ctx, campaign.ID, contact.ID); xerr != nil {
if campaign.OrganizationID == nil {
return nil
}
if xerr := s.advanced.Unsubscribe(ctx, *campaign.OrganizationID, campaign.ID, contact.ID); xerr != nil {
return xerr
}
return nil
+52 -33
View File
@@ -18,6 +18,7 @@ import {
TopbarAction,
} from "@/components/layout/Page";
import { useConfirm } from "@/hooks/context/confirm";
import { usePermission } from "@/hooks/usePermission";
import { useAutomations } from "@/lib/api/hooks/app/automations/useAutomations";
import {
useCreateAutomation,
@@ -50,6 +51,8 @@ export default function AutomationsPage() {
const navigate = useNavigate();
const { data, isLoading } = useAutomations();
const create = useCreateAutomation();
// Building automations is a settings action; others can open and test them.
const canManage = usePermission("MANAGE_SETTINGS");
const automations = data?.automations ?? [];
const enabledCount = automations.filter((a) => a.enabled).length;
@@ -83,12 +86,14 @@ export default function AutomationsPage() {
return (
<Page>
<PageTopbar eyebrow="Automations" subtitle="When something happens in Warmbly, do this across your integrations">
<TopbarAction
onClick={newAutomation}
icon={create.isPending ? <Loader2Icon className="w-3.5 h-3.5 animate-spin" /> : <PlusIcon className="w-3.5 h-3.5" />}
>
New automation
</TopbarAction>
{canManage && (
<TopbarAction
onClick={newAutomation}
icon={create.isPending ? <Loader2Icon className="w-3.5 h-3.5 animate-spin" /> : <PlusIcon className="w-3.5 h-3.5" />}
>
New automation
</TopbarAction>
)}
</PageTopbar>
<StatStrip cols={3}>
@@ -108,23 +113,33 @@ export default function AutomationsPage() {
title="No automations yet"
body="Build a flow: pick a trigger like “meeting booked”, then connect what should happen — ping Slack, create a deal, push to your CRM, or fire a webhook."
cta={
<TopbarAction onClick={newAutomation} icon={<PlusIcon className="w-3.5 h-3.5" />}>
New automation
</TopbarAction>
canManage ? (
<TopbarAction onClick={newAutomation} icon={<PlusIcon className="w-3.5 h-3.5" />}>
New automation
</TopbarAction>
) : undefined
}
/>
) : (
<div className="grid sm:grid-cols-2 lg:grid-cols-3 gap-px bg-slate-200/60 border-b border-slate-200/60">
{automations.map((a, i) => (
<AutomationCard key={a.id} index={i} automation={a} onOpen={() => navigate(`/app/automations/${a.id}`)} />
<AutomationCard
key={a.id}
index={i}
automation={a}
canManage={canManage}
onOpen={() => navigate(`/app/automations/${a.id}`)}
/>
))}
</div>
)}
<div className="mt-8">
<SectionBar label="Start from a template" />
<TemplateGallery onPick={newFromTemplate} busy={create.isPending} />
</div>
{canManage && (
<div className="mt-8">
<SectionBar label="Start from a template" />
<TemplateGallery onPick={newFromTemplate} busy={create.isPending} />
</div>
)}
</PageBody>
</Page>
);
@@ -160,10 +175,12 @@ function AutomationCard({
automation,
onOpen,
index,
canManage,
}: {
automation: Automation;
onOpen: () => void;
index: number;
canManage: boolean;
}) {
const confirm = useConfirm();
const del = useDeleteAutomation();
@@ -253,25 +270,27 @@ function AutomationCard({
</div>
</div>
<div className="mt-auto pt-3 flex items-center justify-between gap-2">
<span
role="button"
tabIndex={-1}
onClick={toggle}
className="text-[11px] text-slate-500 hover:text-slate-900 underline decoration-dotted underline-offset-2"
>
{a.enabled ? "Turn off" : "Turn on"}
</span>
<span
role="button"
tabIndex={-1}
onClick={remove}
className="opacity-100 md:opacity-0 md:group-hover:opacity-100 h-6 w-6 rounded inline-flex items-center justify-center text-slate-400 hover:text-red-600 hover:bg-red-50 transition-all"
title="Delete automation"
>
<Trash2Icon className="w-3.5 h-3.5" />
</span>
</div>
{canManage && (
<div className="mt-auto pt-3 flex items-center justify-between gap-2">
<span
role="button"
tabIndex={-1}
onClick={toggle}
className="text-[11px] text-slate-500 hover:text-slate-900 underline decoration-dotted underline-offset-2"
>
{a.enabled ? "Turn off" : "Turn on"}
</span>
<span
role="button"
tabIndex={-1}
onClick={remove}
className="opacity-100 md:opacity-0 md:group-hover:opacity-100 h-6 w-6 rounded inline-flex items-center justify-center text-slate-400 hover:text-red-600 hover:bg-red-50 transition-all"
title="Delete automation"
>
<Trash2Icon className="w-3.5 h-3.5" />
</span>
</div>
)}
</motion.button>
);
}
@@ -645,8 +645,8 @@ export default function AutomationFlow({
// Collaboration: claim this automation while the builder is open so a
// teammate sees who's here. Editors show as "editing"; members without the
// integration permission (view-only) show as "viewing".
const canEditAutomation = usePermission("USE_INTEGRATIONS");
// settings permission (view-only) show as "viewing".
const canEditAutomation = usePermission("MANAGE_SETTINGS");
usePresenceResource(`automation:${automation.id}`, canEditAutomation ? "editing" : "viewing");
// This canvas has its own flow-space cursor layer; silence the page layer.
useSuppressGlobalCursors();
@@ -694,13 +694,14 @@ export default function AutomationFlow({
// only the ones a hand touched). Silent on the server: no audit, no
// updated_at bump, so it never nudges a teammate's editor.
const persistLayout = React.useCallback(() => {
if (!canEditAutomation) return;
const positions = nodesRef.current.map((n) => ({
id: n.id,
x: Math.round(n.position.x),
y: Math.round(n.position.y),
}));
if (positions.length) layoutMutateRef.current({ id: automation.id, positions });
}, [automation.id]);
}, [automation.id, canEditAutomation]);
// Debounced so a flurry of drags coalesces into one write.
const commitLayout = React.useCallback(() => {
@@ -1267,7 +1268,8 @@ export default function AutomationFlow({
// action steps the user toggled off. Persists the canvas first only when there
// are unsaved edits, so we test what's on screen.
const runTest = async (data?: Record<string, unknown>, skipNodeIds?: string[]) => {
if (dirty && !(await save())) return;
// A view-only member tests the saved flow; only an editor can persist first.
if (dirty && canEditAutomation && !(await save())) return;
try {
const res = await test.mutateAsync({ id: automation.id, data, skipNodeIds });
setTestResult(res);
@@ -1324,9 +1326,10 @@ export default function AutomationFlow({
role="switch"
aria-checked={enabled}
aria-label="Enable automation"
disabled={!canEditAutomation}
onClick={() => setEnabled((v) => !v)}
title={enabled ? "Automation is live" : "Automation is paused"}
className="inline-flex h-7 cursor-pointer select-none items-center gap-2 rounded-md outline-none focus-visible:ring-2 focus-visible:ring-sky-200"
className="inline-flex h-7 cursor-pointer select-none items-center gap-2 rounded-md outline-none focus-visible:ring-2 focus-visible:ring-sky-200 disabled:cursor-default disabled:opacity-70"
>
<span
className={cn(
@@ -1391,7 +1394,7 @@ export default function AutomationFlow({
<span className="hidden md:inline">Tidy up</span>
</button>
<PermissionButton
permission="USE_INTEGRATIONS"
permission="MANAGE_SETTINGS"
type="button"
onClick={save}
disabled={!dirty || update.isPending}