mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-06 00:02:15 +00:00
feat: escape Slack's reserved characters in every text field of the automation Slack card, quote backslashes and double quotes in the Close lead lookup, and route Salesforce automation and push writes only through the native Salesforce sync by removing the unused direct SOQL fallback
This commit is contained in:
@@ -93,7 +93,7 @@ Each action node picks a **Run** target (an integration, or Warmbly built-in) an
|
||||
| Send a webhook | HTTP request to a URL you provide |
|
||||
| Create / update HubSpot contact, Pipedrive person, Salesforce record, Close lead | Upserts the record. In [HubSpot mode](/guides/hubspot/) contacts already sync on their own, so the HubSpot action is only needed for fields the sync does not map. Salesforce finds the Lead or Contact, or creates one, with the [connection's rules](/guides/salesforce/#how-people-are-matched) |
|
||||
|
||||
Slack, Discord, and webhook actions take an optional message template. Slack and Discord arrive as a branded card in Warmbly's accent color with contact and subject fields, not a plain line.
|
||||
Slack, Discord, and webhook actions take an optional message template. Slack and Discord arrive as a branded card in Warmbly's accent color with contact and subject fields, not a plain line. Slack shows the text as written: `<`, `>` and `&` are escaped, so a subject or template cannot form a Slack link or mention.
|
||||
|
||||
| Built-in action | What it does |
|
||||
| --- | --- |
|
||||
|
||||
@@ -108,6 +108,12 @@ func truncateRunes(s string, max int) string {
|
||||
return string(r[:max-1]) + "…"
|
||||
}
|
||||
|
||||
// slackEscape escapes the characters Slack reserves for links and mentions, so
|
||||
// text from a contact or a reply renders as written.
|
||||
func slackEscape(s string) string {
|
||||
return strings.NewReplacer("&", "&", "<", "<", ">", ">").Replace(s)
|
||||
}
|
||||
|
||||
// notifyFields turns the message's contact/subject into structured key-value
|
||||
// fields shared by the Slack and Discord cards (omitted when empty).
|
||||
func (m eventMessage) notifyFields() (contact, subject string) {
|
||||
@@ -125,28 +131,28 @@ func slackPostMessage(ctx context.Context, token, channel string, msg eventMessa
|
||||
}
|
||||
attachment := map[string]any{
|
||||
"color": notifyAccentHex,
|
||||
"fallback": msg.plainText(),
|
||||
"title": truncateRunes(msg.Title, 256),
|
||||
"fallback": slackEscape(msg.plainText()),
|
||||
"title": slackEscape(truncateRunes(msg.Title, 256)),
|
||||
"footer": "Warmbly",
|
||||
"ts": time.Now().Unix(),
|
||||
}
|
||||
if msg.Custom != "" {
|
||||
attachment["text"] = truncateRunes(msg.Custom, 3000)
|
||||
attachment["text"] = slackEscape(truncateRunes(msg.Custom, 3000))
|
||||
}
|
||||
contact, subject := msg.notifyFields()
|
||||
var fields []map[string]any
|
||||
if contact != "" {
|
||||
fields = append(fields, map[string]any{"title": "Contact", "value": contact, "short": true})
|
||||
fields = append(fields, map[string]any{"title": "Contact", "value": slackEscape(contact), "short": true})
|
||||
}
|
||||
if subject != "" {
|
||||
fields = append(fields, map[string]any{"title": "Subject", "value": subject, "short": true})
|
||||
fields = append(fields, map[string]any{"title": "Subject", "value": slackEscape(subject), "short": true})
|
||||
}
|
||||
if len(fields) > 0 {
|
||||
attachment["fields"] = fields
|
||||
}
|
||||
body, _ := json.Marshal(map[string]any{
|
||||
"channel": channel,
|
||||
"text": msg.Title,
|
||||
"text": slackEscape(msg.Title),
|
||||
"attachments": []map[string]any{attachment},
|
||||
})
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, "https://slack.com/api/chat.postMessage", bytes.NewReader(body))
|
||||
@@ -411,7 +417,7 @@ func closeUpsertLead(ctx context.Context, apiKey, email string, props map[string
|
||||
|
||||
// Idempotency guard: skip if a lead already has this email address.
|
||||
searchURL := "https://api.close.com/api/v1/lead/?_fields=id&query=" +
|
||||
url.QueryEscape("email_address:\""+email+"\"")
|
||||
url.QueryEscape("email_address:\""+strings.NewReplacer(`\`, `\\`, `"`, `\"`).Replace(email)+"\"")
|
||||
var search struct {
|
||||
Data []struct {
|
||||
ID string `json:"id"`
|
||||
@@ -474,78 +480,3 @@ func closeJSON(ctx context.Context, method, reqURL, apiKey string, body []byte,
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// salesforceAPIVersion is the REST API version actions target. Salesforce keeps
|
||||
// old versions live for years, so pinning one keeps request shapes stable.
|
||||
const salesforceAPIVersion = "v59.0"
|
||||
|
||||
// salesforceUpsertContact creates or updates a Salesforce Contact keyed by email
|
||||
// using the caller-projected props. instanceURL is the connected org's API host,
|
||||
// captured at OAuth time and stored in the connection's display fields. LastName
|
||||
// is mandatory on the Contact object, so we fall back to the email when unset.
|
||||
func salesforceUpsertContact(ctx context.Context, token, instanceURL, email string, props map[string]any) error {
|
||||
instanceURL, err := SalesforceInstanceURL(instanceURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("salesforce instance url unavailable; reconnect the integration")
|
||||
}
|
||||
if email == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
fields := map[string]any{}
|
||||
for k, v := range props {
|
||||
fields[k] = v
|
||||
}
|
||||
fields["Email"] = email
|
||||
if strProp(fields, "LastName") == "" {
|
||||
fields["LastName"] = email // LastName is a required Contact field.
|
||||
}
|
||||
|
||||
// Find an existing contact by email (SOQL — escape embedded single quotes).
|
||||
soql := "SELECT Id FROM Contact WHERE Email = '" + strings.ReplaceAll(email, "'", "\\'") + "' LIMIT 1"
|
||||
queryURL := instanceURL + "/services/data/" + salesforceAPIVersion + "/query?q=" + url.QueryEscape(soql)
|
||||
var q struct {
|
||||
Records []struct {
|
||||
ID string `json:"Id"`
|
||||
} `json:"records"`
|
||||
}
|
||||
if err := salesforceJSON(ctx, http.MethodGet, queryURL, token, nil, &q); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
base := instanceURL + "/services/data/" + salesforceAPIVersion + "/sobjects/Contact"
|
||||
body, _ := json.Marshal(fields)
|
||||
if len(q.Records) > 0 {
|
||||
// PATCH returns 204 No Content on success.
|
||||
return salesforceJSON(ctx, http.MethodPatch, base+"/"+q.Records[0].ID, token, body, nil)
|
||||
}
|
||||
return salesforceJSON(ctx, http.MethodPost, base+"/", token, body, nil)
|
||||
}
|
||||
|
||||
func salesforceJSON(ctx context.Context, method, reqURL, token string, body []byte, dst any) error {
|
||||
var reader io.Reader
|
||||
if body != nil {
|
||||
reader = bytes.NewReader(body)
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, method, reqURL, reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
req.Header.Set("Authorization", "Bearer "+token)
|
||||
if body != nil {
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
}
|
||||
resp, err := actionHTTP.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
return fmt.Errorf("salesforce %s: HTTP %d", method, resp.StatusCode)
|
||||
}
|
||||
if dst != nil && len(raw) > 0 {
|
||||
return json.Unmarshal(raw, dst)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -189,30 +189,25 @@ func (s *service) execAction(ctx context.Context, target repository.DispatchTarg
|
||||
return pipedriveUpsertPerson(ctx, token, contactEmail(data), props)
|
||||
|
||||
case models.IntegrationActionSalesforceUpsert:
|
||||
if s.salesforce != nil {
|
||||
ev := map[string]any{}
|
||||
for k, v := range data {
|
||||
ev[k] = v
|
||||
}
|
||||
// An automation's own field map overrides the connection's rules.
|
||||
if len(autoCfg.FieldMap) > 0 {
|
||||
ev["_salesforce_fields"] = projectFields(autoCfg.FieldMap, eventSource(data))
|
||||
}
|
||||
if err := s.salesforce.UpsertFromEvent(ctx, sub.OrganizationID, sub.ConnectionID, ev); err != nil {
|
||||
if errors.Is(err, ErrPushReauth) {
|
||||
return errReauthRequired
|
||||
}
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
// Salesforce writes go through the native sync only.
|
||||
if s.salesforce == nil {
|
||||
return errors.New("salesforce sync is not available on this process")
|
||||
}
|
||||
token, terr := s.accessTokenFor(ctx, &target.Secrets)
|
||||
if terr != nil {
|
||||
return errReauthRequired
|
||||
ev := map[string]any{}
|
||||
for k, v := range data {
|
||||
ev[k] = v
|
||||
}
|
||||
instanceURL := configString(target.Secrets.Conn.DisplayFields, "instance_url")
|
||||
props := s.crmProps(ctx, sub, models.IntegrationSalesforce, data, autoCfg)
|
||||
return salesforceUpsertContact(ctx, token, instanceURL, contactEmail(data), props)
|
||||
// An automation's own field map overrides the connection's rules.
|
||||
if len(autoCfg.FieldMap) > 0 {
|
||||
ev["_salesforce_fields"] = projectFields(autoCfg.FieldMap, eventSource(data))
|
||||
}
|
||||
if err := s.salesforce.UpsertFromEvent(ctx, sub.OrganizationID, sub.ConnectionID, ev); err != nil {
|
||||
if errors.Is(err, ErrPushReauth) {
|
||||
return errReauthRequired
|
||||
}
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
|
||||
case models.IntegrationActionCloseUpsert:
|
||||
apiKey := stringFromMap(secretCfg, "api_key", "api_token")
|
||||
|
||||
@@ -85,7 +85,7 @@ func (s *service) PushContacts(ctx context.Context, orgID, connID uuid.UUID, con
|
||||
}
|
||||
|
||||
// Resolve auth once for the whole batch.
|
||||
var token, apiKey, instanceURL string
|
||||
var token, apiKey string
|
||||
switch conn.Provider {
|
||||
case models.IntegrationClose:
|
||||
cfg, cerr := s.openConfig(ctx, sec)
|
||||
@@ -96,15 +96,12 @@ func (s *service) PushContacts(ctx context.Context, orgID, connID uuid.UUID, con
|
||||
if apiKey == "" {
|
||||
return nil, errors.New("no close api key configured")
|
||||
}
|
||||
default: // OAuth CRMs: hubspot, pipedrive, salesforce
|
||||
default: // OAuth CRMs: hubspot, pipedrive
|
||||
tok, terr := s.accessTokenFor(ctx, sec)
|
||||
if terr != nil {
|
||||
return nil, ErrPushReauth
|
||||
}
|
||||
token = tok
|
||||
if conn.Provider == models.IntegrationSalesforce {
|
||||
instanceURL = configString(sec.Conn.DisplayFields, "instance_url")
|
||||
}
|
||||
}
|
||||
|
||||
// Resolve the connection's effective field map once for the whole batch — it
|
||||
@@ -140,8 +137,6 @@ func (s *service) PushContacts(ctx context.Context, orgID, connID uuid.UUID, con
|
||||
aerr = hubspotUpsertContact(ctx, token, ct.Email, props, "Synced from Warmbly")
|
||||
case models.IntegrationPipedrive:
|
||||
aerr = pipedriveUpsertPerson(ctx, token, ct.Email, props)
|
||||
case models.IntegrationSalesforce:
|
||||
aerr = salesforceUpsertContact(ctx, token, instanceURL, ct.Email, props)
|
||||
case models.IntegrationClose:
|
||||
aerr = closeUpsertLead(ctx, apiKey, ct.Email, props)
|
||||
default:
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
package integration
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestSlackEscapeKeepsLinksAndMentionsAsText(t *testing.T) {
|
||||
in := "Re: <https://evil.example|Reset your password> & <!channel>"
|
||||
want := "Re: <https://evil.example|Reset your password> & <!channel>"
|
||||
if got := slackEscape(in); got != want {
|
||||
t.Errorf("slackEscape = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user