From f0ec39febdef8ac96ef5b0ab0346d378875cd57e Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Sun, 4 Oct 2026 02:46:49 -0700 Subject: [PATCH] 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 --- docs/content/docs/api/endpoints.mdx | 5 +- docs/content/docs/api/error-codes.mdx | 6 ++ .../docs/api/reference/integrations.mdx | 8 +- docs/content/docs/guides/team-roles.mdx | 4 +- docs/public/openapi.json | 2 +- internal/api/handler/crm.go | 14 ++- internal/api/handler/integration.go | 14 ++- internal/api/middleware/apikey.go | 8 +- internal/api/middleware/apikey_test.go | 26 ++++++ internal/api/middleware/organization.go | 3 +- internal/api/routes.go | 11 ++- internal/app/advanced/incoming_reply_test.go | 7 ++ internal/app/advanced/lead_copy_reply_test.go | 4 +- internal/app/advanced/reply_actions.go | 5 +- internal/app/advanced/service.go | 23 +++-- internal/app/aitools/tools_crm.go | 4 +- internal/app/crm/service.go | 30 ++++++- .../app/integration/action_provider_test.go | 42 +++++++++ internal/app/integration/dispatch.go | 7 ++ internal/app/integration/native_actions.go | 4 +- internal/app/integration/service.go | 6 ++ internal/app/nativeactions/adapter.go | 4 +- internal/app/research/service.go | 3 + internal/models/integration_capability.go | 6 ++ internal/repository/pg_contact.go | 7 +- internal/repository/pg_crm.go | 8 +- internal/repository/pg_research.go | 11 ++- internal/repository/pg_segment.go | 7 +- internal/tasks/campaign_task.go | 5 +- web/src/app/app/automations/page.tsx | 85 ++++++++++++------- .../app/automations/AutomationFlow.tsx | 15 ++-- 31 files changed, 286 insertions(+), 98 deletions(-) create mode 100644 internal/app/integration/action_provider_test.go diff --git a/docs/content/docs/api/endpoints.mdx b/docs/content/docs/api/endpoints.mdx index 6086e784c..ff4651714 100644 --- a/docs/content/docs/api/endpoints.mdx +++ b/docs/content/docs/api/endpoints.mdx @@ -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 diff --git a/docs/content/docs/api/error-codes.mdx b/docs/content/docs/api/error-codes.mdx index 58fdc79aa..8b17cbdf7 100644 --- a/docs/content/docs/api/error-codes.mdx +++ b/docs/content/docs/api/error-codes.mdx @@ -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. diff --git a/docs/content/docs/api/reference/integrations.mdx b/docs/content/docs/api/reference/integrations.mdx index f7a3c338c..9382f327e 100644 --- a/docs/content/docs/api/reference/integrations.mdx +++ b/docs/content/docs/api/reference/integrations.mdx @@ -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 | |-----------|-----|------|-------------| diff --git a/docs/content/docs/guides/team-roles.mdx b/docs/content/docs/guides/team-roles.mdx index 598fdddbf..75e8f9448 100644 --- a/docs/content/docs/guides/team-roles.mdx +++ b/docs/content/docs/guides/team-roles.mdx @@ -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 | diff --git a/docs/public/openapi.json b/docs/public/openapi.json index 26297fefe..794ad053a 100644 --- a/docs/public/openapi.json +++ b/docs/public/openapi.json @@ -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" ], diff --git a/internal/api/handler/crm.go b/internal/api/handler/crm.go index 3f1935598..0c37ebd43 100644 --- a/internal/api/handler/crm.go +++ b/internal/api/handler/crm.go @@ -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 diff --git a/internal/api/handler/integration.go b/internal/api/handler/integration.go index 2ab1b1d2b..d06df4d69 100644 --- a/internal/api/handler/integration.go +++ b/internal/api/handler/integration.go @@ -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) diff --git a/internal/api/middleware/apikey.go b/internal/api/middleware/apikey.go index 99f011bf9..1979dcc3f 100644 --- a/internal/api/middleware/apikey.go +++ b/internal/api/middleware/apikey.go @@ -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) diff --git a/internal/api/middleware/apikey_test.go b/internal/api/middleware/apikey_test.go index 341b93b2e..c0d39a53c 100644 --- a/internal/api/middleware/apikey_test.go +++ b/internal/api/middleware/apikey_test.go @@ -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) + } + } +} diff --git a/internal/api/middleware/organization.go b/internal/api/middleware/organization.go index 8d3b93dbf..367ec85fc 100644 --- a/internal/api/middleware/organization.go +++ b/internal/api/middleware/organization.go @@ -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 } diff --git a/internal/api/routes.go b/internal/api/routes.go index eb2e9d0e0..0fdffc179 100644 --- a/internal/api/routes.go +++ b/internal/api/routes.go @@ -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) diff --git a/internal/app/advanced/incoming_reply_test.go b/internal/app/advanced/incoming_reply_test.go index 0dc8a7886..0acdd32db 100644 --- a/internal/app/advanced/incoming_reply_test.go +++ b/internal/app/advanced/incoming_reply_test.go @@ -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 diff --git a/internal/app/advanced/lead_copy_reply_test.go b/internal/app/advanced/lead_copy_reply_test.go index 4415bce26..bd337fd9d 100644 --- a/internal/app/advanced/lead_copy_reply_test.go +++ b/internal/app/advanced/lead_copy_reply_test.go @@ -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) diff --git a/internal/app/advanced/reply_actions.go b/internal/app/advanced/reply_actions.go index 275659208..90dabc536 100644 --- a/internal/app/advanced/reply_actions.go +++ b/internal/app/advanced/reply_actions.go @@ -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": diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go index a84ad0dd7..7ae98f5ef 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -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)", diff --git a/internal/app/aitools/tools_crm.go b/internal/app/aitools/tools_crm.go index cf2743734..9d74c5f02 100644 --- a/internal/app/aitools/tools_crm.go +++ b/internal/app/aitools/tools_crm.go @@ -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) diff --git a/internal/app/crm/service.go b/internal/app/crm/service.go index 5ff639a57..c85c8aef8 100644 --- a/internal/app/crm/service.go +++ b/internal/app/crm/service.go @@ -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 diff --git a/internal/app/integration/action_provider_test.go b/internal/app/integration/action_provider_test.go new file mode 100644 index 000000000..e6c441142 --- /dev/null +++ b/internal/app/integration/action_provider_test.go @@ -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) + } +} diff --git a/internal/app/integration/dispatch.go b/internal/app/integration/dispatch.go index 368350cd7..cbc94db4b 100644 --- a/internal/app/integration/dispatch.go +++ b/internal/app/integration/dispatch.go @@ -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) diff --git a/internal/app/integration/native_actions.go b/internal/app/integration/native_actions.go index 949458952..d071c98a6 100644 --- a/internal/app/integration/native_actions.go +++ b/internal/app/integration/native_actions.go @@ -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) diff --git a/internal/app/integration/service.go b/internal/app/integration/service.go index 4a92ef4cb..38223f01d 100644 --- a/internal/app/integration/service.go +++ b/internal/app/integration/service.go @@ -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) diff --git a/internal/app/nativeactions/adapter.go b/internal/app/nativeactions/adapter.go index 1a27e3867..f9781b025 100644 --- a/internal/app/nativeactions/adapter.go +++ b/internal/app/nativeactions/adapter.go @@ -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 diff --git a/internal/app/research/service.go b/internal/app/research/service.go index 7a9681f81..b163b31b1 100644 --- a/internal/app/research/service.go +++ b/internal/app/research/service.go @@ -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") } diff --git a/internal/models/integration_capability.go b/internal/models/integration_capability.go index 3dc960c10..65ca65689 100644 --- a/internal/models/integration_capability.go +++ b/internal/models/integration_capability.go @@ -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 { diff --git a/internal/repository/pg_contact.go b/internal/repository/pg_contact.go index a18cfa26f..3d103fd56 100644 --- a/internal/repository/pg_contact.go +++ b/internal/repository/pg_contact.go @@ -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() } diff --git a/internal/repository/pg_crm.go b/internal/repository/pg_crm.go index 9a5f589ce..defe90641 100644 --- a/internal/repository/pg_crm.go +++ b/internal/repository/pg_crm.go @@ -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++ } diff --git a/internal/repository/pg_research.go b/internal/repository/pg_research.go index 4d7f76b2c..f9ec9d0d9 100644 --- a/internal/repository/pg_research.go +++ b/internal/repository/pg_research.go @@ -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 } diff --git a/internal/repository/pg_segment.go b/internal/repository/pg_segment.go index 51ed90639..39512bd33 100644 --- a/internal/repository/pg_segment.go +++ b/internal/repository/pg_segment.go @@ -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 diff --git a/internal/tasks/campaign_task.go b/internal/tasks/campaign_task.go index 613376abb..f87376373 100644 --- a/internal/tasks/campaign_task.go +++ b/internal/tasks/campaign_task.go @@ -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 diff --git a/web/src/app/app/automations/page.tsx b/web/src/app/app/automations/page.tsx index 4f3dc352a..bf8881e16 100644 --- a/web/src/app/app/automations/page.tsx +++ b/web/src/app/app/automations/page.tsx @@ -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 ( - : } - > - New automation - + {canManage && ( + : } + > + New automation + + )} @@ -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={ - }> - New automation - + canManage ? ( + }> + New automation + + ) : undefined } /> ) : (
{automations.map((a, i) => ( - navigate(`/app/automations/${a.id}`)} /> + navigate(`/app/automations/${a.id}`)} + /> ))}
)} -
- - -
+ {canManage && ( +
+ + +
+ )}
); @@ -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({ -
- - {a.enabled ? "Turn off" : "Turn on"} - - - - -
+ {canManage && ( +
+ + {a.enabled ? "Turn off" : "Turn on"} + + + + +
+ )} ); } diff --git a/web/src/components/app/automations/AutomationFlow.tsx b/web/src/components/app/automations/AutomationFlow.tsx index 6b50fb274..10f5d2da0 100644 --- a/web/src/components/app/automations/AutomationFlow.tsx +++ b/web/src/components/app/automations/AutomationFlow.tsx @@ -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, 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" > Tidy up