Add slack integration

Signed-off-by: Sacha Al Himdani <sacha@getprobo.com>
This commit is contained in:
Sacha Al Himdani
2025-10-15 23:41:59 +02:00
parent 8302b11614
commit de004ce8d7
38 changed files with 2621 additions and 1578 deletions

View File

@@ -27,8 +27,8 @@ type (
ProtocolType string
Connector interface {
Initiate(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (string, error)
Complete(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (Connection, error)
Initiate(ctx context.Context, provider string, organizationID gid.GID, r *http.Request) (string, error)
Complete(ctx context.Context, r *http.Request) (Connection, *gid.GID, error)
}
Connection interface {
@@ -41,20 +41,28 @@ type (
)
const (
ProtocolOAuth2 ProtocolType = "oauth2"
ProtocolOAuth2 ProtocolType = "OAUTH2"
)
func UnmarshalConnection(prtcl ProtocolType, data []byte) (Connection, error) {
func UnmarshalConnection(protocol string, provider string, data []byte) (Connection, error) {
switch protocol {
case string(ProtocolOAuth2):
switch provider {
case SlackProvider:
var slackConn SlackConnection
if err := json.Unmarshal(data, &slackConn); err != nil {
return nil, fmt.Errorf("cannot unmarshal slack connection: %w", err)
}
return &slackConn, nil
switch prtcl {
case ProtocolOAuth2:
var conn OAuth2Connection
if err := json.Unmarshal(data, &conn); err != nil {
return nil, fmt.Errorf("cannot unmarshal oauth2 connection: %w", err)
default:
var conn OAuth2Connection
if err := json.Unmarshal(data, &conn); err != nil {
return nil, fmt.Errorf("cannot unmarshal oauth2 connection: %w", err)
}
return &conn, nil
}
return &conn, nil
}
return nil, fmt.Errorf("unknown connection type: %s", prtcl)
return nil, fmt.Errorf("unknown connection protocol: %s", protocol)
}

View File

@@ -15,9 +15,11 @@
package connector
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"strings"
@@ -46,7 +48,7 @@ type (
OAuth2State struct {
OrganizationID string `json:"oid"`
ConnectorID string `json:"cid"`
Provider string `json:"provider"`
}
OAuth2Connection struct {
@@ -66,29 +68,33 @@ var (
OAuth2TokenTTL = 10 * time.Minute
)
func (c *OAuth2Connector) Initiate(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (string, error) {
stateData := OAuth2State{OrganizationID: organizationID.String(), ConnectorID: connectorID}
func (c *OAuth2Connector) Initiate(ctx context.Context, provider string, organizationID gid.GID, r *http.Request) (string, error) {
stateData := OAuth2State{
OrganizationID: organizationID.String(),
Provider: provider,
}
state, err := statelesstoken.NewToken(c.ClientSecret, OAuth2TokenType, OAuth2TokenTTL, stateData)
if err != nil {
return "", fmt.Errorf("cannot create state token: %w", err)
}
redirectURI, err := url.Parse(c.RedirectURI)
redirectURI := c.RedirectURI
redirectURIParsed, err := url.Parse(redirectURI)
if err != nil {
return "", fmt.Errorf("cannot parse redirect URI: %w", err)
}
redirectQuery := url.Values{}
redirectQuery.Set("organization_id", organizationID.String())
redirectQuery.Set("connector_id", connectorID)
redirectQuery.Set("continue", r.URL.Query().Get("continue"))
redirectURI.RawQuery = redirectQuery.Encode()
q := redirectURIParsed.Query()
q.Set("provider", provider)
if continueURL := r.URL.Query().Get("continue"); continueURL != "" {
q.Set("continue", continueURL)
}
redirectURIParsed.RawQuery = q.Encode()
redirectURI = redirectURIParsed.String()
authCodeQuery := url.Values{}
authCodeQuery.Set("state", state)
authCodeQuery.Set("client_id", c.ClientID)
authCodeQuery.Set("redirect_uri", redirectURI.String())
authCodeQuery.Set("redirect_uri", redirectURI)
authCodeQuery.Set("response_type", "code")
authCodeQuery.Set("scope", strings.Join(c.Scopes, " "))
@@ -102,51 +108,59 @@ func (c *OAuth2Connector) Initiate(ctx context.Context, connectorID string, orga
return u.String(), nil
}
func (c *OAuth2Connector) Complete(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (Connection, error) {
func (c *OAuth2Connector) Complete(ctx context.Context, r *http.Request) (Connection, *gid.GID, error) {
provider := r.URL.Query().Get("provider")
if provider == "" {
return nil, nil, fmt.Errorf("missing provider in query parameters")
}
code := r.URL.Query().Get("code")
if code == "" {
return nil, fmt.Errorf("no code in request")
return nil, nil, fmt.Errorf("no code in request")
}
state := r.URL.Query().Get("state")
if state == "" {
return nil, fmt.Errorf("no state in request")
stateToken := r.URL.Query().Get("state")
if stateToken == "" {
return nil, nil, fmt.Errorf("no state in request")
}
payload, err := statelesstoken.ValidateToken[OAuth2State](c.ClientSecret, OAuth2TokenType, state)
payload, err := statelesstoken.ValidateToken[OAuth2State](c.ClientSecret, OAuth2TokenType, stateToken)
if err != nil {
return nil, fmt.Errorf("cannot validate state token: %w", err)
return nil, nil, fmt.Errorf("cannot validate state token: %w", err)
}
if payload.Data.OrganizationID != organizationID.String() {
return nil, fmt.Errorf("invalid organization ID")
if payload.Data.Provider != provider {
return nil, nil, fmt.Errorf("provider mismatch: state has %q, query has %q", payload.Data.Provider, provider)
}
if payload.Data.ConnectorID != connectorID {
return nil, fmt.Errorf("invalid connector ID")
}
redirectURI, err := url.Parse(c.RedirectURI)
organizationID, err := gid.ParseGID(payload.Data.OrganizationID)
if err != nil {
return nil, fmt.Errorf("cannot parse redirect URI: %w", err)
return nil, nil, fmt.Errorf("cannot parse organization ID: %w", err)
}
redirectQuery := url.Values{}
redirectQuery.Set("organization_id", organizationID.String())
redirectQuery.Set("connector_id", connectorID)
redirectURI.RawQuery = redirectQuery.Encode()
redirectURI := c.RedirectURI
redirectURIParsed, err := url.Parse(redirectURI)
if err != nil {
return nil, nil, fmt.Errorf("cannot parse redirect URI: %w", err)
}
q := redirectURIParsed.Query()
q.Set("provider", provider)
if continueURL := r.URL.Query().Get("continue"); continueURL != "" {
q.Set("continue", continueURL)
}
redirectURIParsed.RawQuery = q.Encode()
redirectURI = redirectURIParsed.String()
tokenRequestData := url.Values{}
tokenRequestData.Set("client_id", c.ClientID)
tokenRequestData.Set("client_secret", c.ClientSecret)
tokenRequestData.Set("code", code)
tokenRequestData.Set("redirect_uri", redirectURI.String())
tokenRequestData.Set("redirect_uri", redirectURI)
tokenRequestData.Set("grant_type", "authorization_code")
tokenRequest, err := http.NewRequestWithContext(ctx, "POST", c.TokenURL, strings.NewReader(tokenRequestData.Encode()))
tokenRequest, err := http.NewRequestWithContext(ctx, http.MethodPost, c.TokenURL, strings.NewReader(tokenRequestData.Encode()))
if err != nil {
return nil, fmt.Errorf("cannot create token request: %w", err)
return nil, nil, fmt.Errorf("cannot create token request: %w", err)
}
tokenRequest.Header.Set("Content-Type", "application/x-www-form-urlencoded; charset=utf-8")
@@ -155,35 +169,32 @@ func (c *OAuth2Connector) Complete(ctx context.Context, connectorID string, orga
tokenResp, err := http.DefaultClient.Do(tokenRequest)
if err != nil {
return nil, fmt.Errorf("cannot post token URL: %w", err)
return nil, nil, fmt.Errorf("cannot post token URL: %w", err)
}
defer tokenResp.Body.Close()
if tokenResp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("token response status: %d", tokenResp.StatusCode)
return nil, nil, fmt.Errorf("token response status: %d", tokenResp.StatusCode)
}
type tokenResponse struct {
AccessToken string `json:"access_token"`
RefreshToken string `json:"refresh_token"`
ExpiresAt time.Time `json:"expires_at"`
Scope string `json:"scope"`
TokenType string `json:"token_type"`
}
var token tokenResponse
err = json.NewDecoder(tokenResp.Body).Decode(&token)
body, err := io.ReadAll(tokenResp.Body)
if err != nil {
return nil, fmt.Errorf("cannot decode token response: %w", err)
return nil, nil, fmt.Errorf("cannot read token response body: %w", err)
}
return &OAuth2Connection{
AccessToken: token.AccessToken,
RefreshToken: token.RefreshToken,
ExpiresAt: token.ExpiresAt,
Scope: token.Scope,
TokenType: token.TokenType,
}, nil
var oauth2Conn OAuth2Connection
var buf bytes.Buffer
buf.Write(body)
err = json.NewDecoder(&buf).Decode(&oauth2Conn)
if err != nil {
return nil, nil, fmt.Errorf("cannot decode token response: %w", err)
}
if provider == SlackProvider {
return ParseSlackTokenResponse(body, oauth2Conn, organizationID)
}
return &oauth2Conn, &organizationID, nil
}
func (c *OAuth2Connection) Type() ProtocolType {

View File

@@ -36,40 +36,40 @@ func NewConnectorRegistry() *ConnectorRegistry {
}
}
func (cr *ConnectorRegistry) Register(connectorID string, connector Connector) error {
func (cr *ConnectorRegistry) Register(provider string, connector Connector) error {
cr.Lock()
defer cr.Unlock()
if _, ok := cr.connectors[connectorID]; ok {
return fmt.Errorf("connector %q already registered", connectorID)
if _, ok := cr.connectors[provider]; ok {
return fmt.Errorf("connector %q already registered", provider)
}
cr.connectors[connectorID] = connector
cr.connectors[provider] = connector
return nil
}
func (cr *ConnectorRegistry) Get(connectorID string) (Connector, error) {
func (cr *ConnectorRegistry) Get(provider string) (Connector, error) {
cr.RLock()
defer cr.RUnlock()
connector, ok := cr.connectors[connectorID]
connector, ok := cr.connectors[provider]
if !ok {
return nil, fmt.Errorf("connector %q not found", connectorID)
return nil, fmt.Errorf("connector %q not found", provider)
}
return connector, nil
}
func (cr *ConnectorRegistry) Initiate(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (string, error) {
connector, err := cr.Get(connectorID)
func (cr *ConnectorRegistry) Initiate(ctx context.Context, provider string, organizationID gid.GID, r *http.Request) (string, error) {
connector, err := cr.Get(provider)
if err != nil {
return "", fmt.Errorf("cannot initiate connector: %w", err)
}
return connector.Initiate(ctx, connectorID, organizationID, r)
return connector.Initiate(ctx, provider, organizationID, r)
}
func (cr *ConnectorRegistry) Complete(ctx context.Context, connectorID string, organizationID gid.GID, r *http.Request) (Connection, error) {
connector, err := cr.Get(connectorID)
func (cr *ConnectorRegistry) Complete(ctx context.Context, provider string, r *http.Request) (Connection, *gid.GID, error) {
connector, err := cr.Get(provider)
if err != nil {
return nil, fmt.Errorf("cannot complete connector: %w", err)
return nil, nil, fmt.Errorf("cannot complete connector: %w", err)
}
return connector.Complete(ctx, connectorID, organizationID, r)
return connector.Complete(ctx, r)
}

134
pkg/connector/slack.go Normal file
View File

@@ -0,0 +1,134 @@
// Copyright (c) 2025 Probo Inc <hello@getprobo.com>.
//
// Permission to use, copy, modify, and/or distribute this software for any
// purpose with or without fee is hereby granted, provided that the above
// copyright notice and this permission notice appear in all copies.
//
// THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES WITH
// REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF MERCHANTABILITY
// AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR ANY SPECIAL, DIRECT,
// INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES WHATSOEVER RESULTING FROM
// LOSS OF USE, DATA OR PROFITS, WHETHER IN AN ACTION OF CONTRACT, NEGLIGENCE OR
// OTHER TORTIOUS ACTION, ARISING OUT OF OR IN CONNECTION WITH THE USE OR
// PERFORMANCE OF THIS SOFTWARE.
package connector
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/getprobo/probo/pkg/gid"
)
type (
SlackConnection struct {
OAuth2Connection
Settings SlackSettings `json:"settings"`
}
SlackSettings struct {
WebhookURL string `json:"webhook_url,omitempty"` // Encrypted
Channel string `json:"channel,omitempty"`
ChannelID string `json:"channel_id,omitempty"`
}
IncomingWebhook struct {
URL string `json:"url"`
Channel string `json:"channel"`
ChannelID string `json:"channel_id"`
}
SlackTokenResponse struct {
IncomingWebhook *IncomingWebhook `json:"incoming_webhook,omitempty"`
}
)
const (
SlackProvider = "SLACK"
)
var _ Connection = (*SlackConnection)(nil)
func (c *SlackConnection) Type() ProtocolType {
return ProtocolOAuth2
}
func (c *SlackConnection) Client(ctx context.Context) (*http.Client, error) {
return c.OAuth2Connection.Client(ctx)
}
func (c SlackConnection) MarshalJSON() ([]byte, error) {
return json.Marshal(&struct {
Type string `json:"type"`
AccessToken string `json:"access_token"`
RefreshToken string `json:"refresh_token,omitempty"`
ExpiresAt time.Time `json:"expires_at"`
TokenType string `json:"token_type"`
Scope string `json:"scope,omitempty"`
WebhookURL string `json:"webhook_url,omitempty"`
}{
Type: string(ProtocolOAuth2),
AccessToken: c.OAuth2Connection.AccessToken,
RefreshToken: c.OAuth2Connection.RefreshToken,
ExpiresAt: c.OAuth2Connection.ExpiresAt,
TokenType: c.OAuth2Connection.TokenType,
Scope: c.OAuth2Connection.Scope,
WebhookURL: c.Settings.WebhookURL,
})
}
func (c *SlackConnection) UnmarshalJSON(data []byte) error {
aux := &struct {
Type string `json:"type"`
AccessToken string `json:"access_token"`
RefreshToken string `json:"refresh_token,omitempty"`
ExpiresAt time.Time `json:"expires_at"`
TokenType string `json:"token_type"`
Scope string `json:"scope,omitempty"`
WebhookURL string `json:"webhook_url,omitempty"`
}{}
if err := json.Unmarshal(data, aux); err != nil {
return err
}
c.OAuth2Connection = OAuth2Connection{
AccessToken: aux.AccessToken,
RefreshToken: aux.RefreshToken,
ExpiresAt: aux.ExpiresAt,
TokenType: aux.TokenType,
Scope: aux.Scope,
}
c.Settings.WebhookURL = aux.WebhookURL
return nil
}
func ParseSlackTokenResponse(body []byte, oauth2Conn OAuth2Connection, organizationID gid.GID) (*SlackConnection, *gid.GID, error) {
var slackResponse SlackTokenResponse
var buf bytes.Buffer
buf.Write(body)
if err := json.NewDecoder(&buf).Decode(&slackResponse); err != nil {
return nil, nil, fmt.Errorf("cannot decode Slack token response: %w", err)
}
if slackResponse.IncomingWebhook == nil {
return nil, nil, fmt.Errorf("incoming webhook is required for Slack")
}
settings := SlackSettings{
WebhookURL: slackResponse.IncomingWebhook.URL,
Channel: slackResponse.IncomingWebhook.Channel,
ChannelID: slackResponse.IncomingWebhook.ChannelID,
}
return &SlackConnection{
OAuth2Connection: oauth2Conn,
Settings: settings,
}, &organizationID, nil
}