Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions backend/danswer/configs/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ class DocumentSource(str, Enum):
CLICKUP = "clickup"
MEDIAWIKI = "mediawiki"
WIKIPEDIA = "wikipedia"
SFKBARTICLES = "sfkbarticles"
S3 = "s3"
R2 = "r2"
GOOGLE_CLOUD_STORAGE = "google_cloud_storage"
Expand Down
2 changes: 2 additions & 0 deletions backend/danswer/connectors/factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
from danswer.connectors.productboard.connector import ProductboardConnector
from danswer.connectors.requesttracker.connector import RequestTrackerConnector
from danswer.connectors.salesforce.connector import SalesforceConnector
from danswer.connectors.sfkbarticles.connector import SfKbArticlesConnector
from danswer.connectors.sharepoint.connector import SharepointConnector
from danswer.connectors.slab.connector import SlabConnector
from danswer.connectors.slack.connector import SlackPollConnector
Expand Down Expand Up @@ -86,6 +87,7 @@ def identify_connector_class(
DocumentSource.SHAREPOINT: SharepointConnector,
DocumentSource.TEAMS: TeamsConnector,
DocumentSource.SALESFORCE: SalesforceConnector,
DocumentSource.SFKBARTICLES: SfKbArticlesConnector,
DocumentSource.DISCOURSE: DiscourseConnector,
DocumentSource.AXERO: AxeroConnector,
DocumentSource.CLICKUP: ClickupConnector,
Expand Down
2 changes: 1 addition & 1 deletion backend/danswer/connectors/salesforce/connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -271,4 +271,4 @@ def retrieve_all_source_ids(self) -> set[str]:
}
)
document_batches = connector.load_from_state()
print(next(document_batches))
print(next(document_batches))
2 changes: 1 addition & 1 deletion backend/danswer/connectors/salesforce/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,4 +63,4 @@ def _json_to_natural_language(data: Union[dict, list], indent: int = 0) -> str:
def extract_dict_text(raw_dict: dict) -> str:
processed_dict = _clean_salesforce_dict(raw_dict)
natural_language_dict = _json_to_natural_language(processed_dict)
return natural_language_dict
return natural_language_dict
Empty file.
253 changes: 253 additions & 0 deletions backend/danswer/connectors/sfkbarticles/connector.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,253 @@
import os
import requests

from datetime import datetime
from datetime import timezone
from typing import Any, Tuple

from danswer.configs.app_configs import INDEX_BATCH_SIZE
from danswer.configs.constants import DocumentSource
from danswer.connectors.cross_connector_utils.miscellaneous_utils import time_str_to_utc
from danswer.connectors.interfaces import GenerateDocumentsOutput
from danswer.connectors.interfaces import LoadConnector
from danswer.connectors.interfaces import PollConnector
from danswer.connectors.interfaces import SecondsSinceUnixEpoch
from danswer.connectors.models import BasicExpertInfo
from danswer.connectors.models import Document
from danswer.connectors.models import Section
from danswer.connectors.salesforce.utils import extract_dict_text
from danswer.utils.logger import setup_logger

ID_PREFIX = "SALESFORCE_"
AUTH_URL = "https://login.salesforce.com/services/oauth2/token"

logger = setup_logger()


class SfKbArticlesConnector(LoadConnector, PollConnector):
def __init__(
self,
batch_size: int = INDEX_BATCH_SIZE,
requested_objects: list[str] = [],
) -> None:
self.batch_size = batch_size
self.product_component_list = (
[obj.strip() for obj in requested_objects[0].split(",")]
if requested_objects
else None
)

def load_credentials(self, credentials: dict[str, Any]) -> dict[str, Any] | None:
self.client_id = credentials['sf_client_id']
self.client_secret = credentials['sf_client_secret']
self.username = credentials['sf_username']
self.password = credentials['sf_password']

self.access_token, self.instance_url = self._get_access_token()

self.headers = {
"Authorization": f"Bearer {self.access_token}",
"Content-Type": "application/json",
}

def _get_access_token(self) -> Tuple[str, str]:
"""
Authenticates with Salesforce and retrieves the access token & instance URL.
"""
payload = {
"grant_type": "password",
"client_id": self.client_id,
"client_secret": self.client_secret,
"username": self.username,
"password": f"{self.password}",
}
response = requests.post(AUTH_URL, data=payload)
if response.status_code != 200:
logger.error(f"Authentication failed: {response.text}")
raise Exception("Failed to authenticate with Salesforce.")

data = response.json()
logger.info("Successfully authenticated with Salesforce.")
return data["access_token"], data["instance_url"]

def _convert_object_instance_to_document(
self, object_dict: dict[str, Any]
) -> Document:

salesforce_id = object_dict["Id"]
danswer_salesforce_id = f"{ID_PREFIX}{salesforce_id}"
extracted_link = f"{self.instance_url}/{salesforce_id}"
extracted_doc_updated_at = time_str_to_utc(object_dict["LastModifiedDate"])
extracted_object_text = extract_dict_text(object_dict)
extracted_semantic_identifier = object_dict.get("Title", "Unknown Object")
extracted_primary_owners = [
BasicExpertInfo(
display_name=self._get_name_from_id(object_dict["LastModifiedById"])
)
]

doc = Document(
id=danswer_salesforce_id,
sections=[Section(link=extracted_link, text=extracted_object_text)],
source=DocumentSource.SFKBARTICLES,
semantic_identifier=extracted_semantic_identifier,
doc_updated_at=extracted_doc_updated_at,
primary_owners=extracted_primary_owners,
metadata={},
)
return doc

def _get_name_from_id(self, id: str) -> str:
"""
Fetches the name of a Salesforce user based on their ID.
"""
query = f"SELECT Name FROM User WHERE Id = '{id}'"
url = f"{self.instance_url}/services/data/v56.0/query"
params = {"q": query}

response = requests.get(url, headers=self.headers, params=params)
if response.status_code != 200:
logger.error(f"Failed to fetch name for ID {id}: {response.text}")
return "Unknown"

data = response.json()
records = data.get("records", [])
if not records:
logger.warning(f"No name found for ID {id}")
return "Unknown"

return records[0].get("Name", "Unknown")

def build_salesforce_query(self,
parent_objects: list[str],
start: datetime | None = None,
end: datetime | None = None) -> str:
"""
Builds the Salesforce query dynamically.
If product_component_list is empty, it fetches all Product_Component__c values.
Otherwise, it filters using the provided product_component_list.
"""
if not parent_objects:
product_filter = "" # No filter, fetch all
else:
product_components_str = ", ".join(
[
f"'{component.strip()}'"
for component in parent_objects
if component.strip()
]
)
product_filter = f"AND Product_Component__c IN ({product_components_str})"

query = f"""
SELECT
Id,
Title,
Summary,
Product_Component__c,
Product_Component_Version__c,
Question_Problem__c,
Resolution__c,
Sub_Component__c,
IsVisibleInPkb,
IsVisibleInCsp,
IsVisibleInPrm,
ArticleCreatedDate,
ArticleNumber,
CreatedDate,
ArticleTotalViewCount,
FirstPublishedDate,
IsDeleted,
IsLatestVersion,
Language,
LastModifiedDate,
LastModifiedById,
LastPublishedDate,
Orchestrator_Version__c,
Studio_Version__c
FROM Knowledge__kav
WHERE Language='en_US'
AND PublishStatus = 'Online'
AND IsDeleted = FALSE
{product_filter}
""".strip()

if start:
start_str = start.strftime("%Y-%m-%dT%H:%M:%S.000Z")
query += f" AND LastModifiedDate >= {start_str}"
if end:
end_str = end.strftime("%Y-%m-%dT%H:%M:%S.000Z")
query += f" AND LastModifiedDate <= {end_str}"

return query

def _fetch_from_salesforce(
self,
start: datetime | None = None,
end: datetime | None = None,
) -> GenerateDocumentsOutput:

query = self.build_salesforce_query(self.product_component_list, start, end)
doc_batch: list[Document] = []
query_results: dict = {}

url = f"{self.instance_url}/services/data/v56.0/query"
params = {"q": query}
query_result = requests.get(
url, headers=self.headers, params=params if "q" in params else None
)

while url:
query_result = requests.get(
url, headers=self.headers, params=params if "q" in params else None
)
data = query_result.json()

if isinstance(data, list):
error_message = "; ".join(error.get("message", "Unknown error") for error in data)
raise Exception(f"Salesforce API error: {error_message}")

if "records" in data:
for record_dict in data["records"]:
query_results.setdefault(record_dict["Id"], {}).update(record_dict)

url = data.get("nextRecordsUrl", None)
if url:
url = f"{self.instance_url}{url}"

for combined_object_dict in query_results.values():
doc_batch.append(
self._convert_object_instance_to_document(combined_object_dict)
)
if len(doc_batch) > self.batch_size:
yield doc_batch
doc_batch = []

yield doc_batch

def load_from_state(self) -> GenerateDocumentsOutput:
return self._fetch_from_salesforce()

def poll_source(
self, start: SecondsSinceUnixEpoch, end: SecondsSinceUnixEpoch
) -> GenerateDocumentsOutput:
start_datetime = datetime.fromtimestamp(start, tz=timezone.utc)
end_datetime = datetime.fromtimestamp(end, tz=timezone.utc)
return self._fetch_from_salesforce(start=start_datetime, end=end_datetime)


if __name__ == "__main__":
connector = SfKbArticlesConnector(
requested_objects=os.environ["REQUESTED_OBJECTS"].split(",")
)

connector.load_credentials(
{
"sf_client_id": os.environ["SF_CLIENT_ID"],
"sf_client_secret": os.environ["SF_CLIENT_SECRET"],
"sf_username": os.environ["SF_USERNAME"],
"sf_password": os.environ["SF_PASSWORD"],
}
)
document_batches = connector.load_from_state()
print(next(document_batches))
94 changes: 94 additions & 0 deletions backend/danswer/connectors/sfkbarticles/utils.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
import re
from bs4 import BeautifulSoup
from typing import Union
from danswer.utils.logger import setup_logger

SF_JSON_FILTER = r"Id$|Date$|stamp$|url$"


def _clean_salesforce_dict(data: Union[dict, list]) -> Union[dict, list]:
if isinstance(data, dict):
if "records" in data.keys():
data = data["records"]
if isinstance(data, dict):
if "attributes" in data.keys():
if isinstance(data["attributes"], dict):
data.update(data.pop("attributes"))

if isinstance(data, dict):
filtered_dict = {}
for key, value in data.items():
if not re.search(SF_JSON_FILTER, key, re.IGNORECASE):
if "__c" in key: # remove the custom object indicator for display
key = key[:-3]
if isinstance(value, (dict, list)):
filtered_value = _clean_salesforce_dict(value)
if filtered_value:
filtered_dict[key] = filtered_value
elif value is not None:
filtered_dict[key] = value
return filtered_dict
elif isinstance(data, list):
filtered_list = []
for item in data:
if isinstance(item, (dict, list)):
filtered_item = _clean_salesforce_dict(item)
if filtered_item:
filtered_list.append(filtered_item)
elif item is not None:
filtered_list.append(filtered_item)
return filtered_list
else:
return data


def _json_to_natural_language(data: Union[dict, list], indent: int = 0) -> str:
result = []
indent_str = " " * indent

if isinstance(data, dict):
for key, value in data.items():
if isinstance(value, (dict, list)):
result.append(f"{indent_str}{key}:")
result.append(_json_to_natural_language(value, indent + 2))
else:
result.append(f"{indent_str}{key}: {value}")
elif isinstance(data, list):
for item in data:
result.append(_json_to_natural_language(item, indent))
else:
result.append(f"{indent_str}{data}")

return "\n".join(result)


def extract_dict_text(raw_dict: dict) -> str:
processed_dict = _clean_salesforce_dict(raw_dict)

if 'Resolution' in processed_dict:
processed_dict['Resolution'] = clean_html(processed_dict['Resolution'])

natural_language_dict = _json_to_natural_language(processed_dict)
return natural_language_dict


def clean_html(html_content: str) -> str:
"""
Cleans the HTML content by removing tags and returning the plain text,
while preserving the links.
:param html_content: HTML content as string
:return: Cleaned text with preserved links
"""
soup = BeautifulSoup(html_content, "lxml")

# Replace <a> tags with their text and the href as a link in brackets
for a_tag in soup.find_all('a', href=True):
a_tag.insert_before(f"[{a_tag.get_text()}]({a_tag['href']})")
a_tag.decompose()

cleaned_text = soup.get_text(separator="\n", strip=True)

# Replace non-breaking spaces (\xa0) with regular spaces
cleaned_text = cleaned_text.replace('\xa0', ' ')

return cleaned_text
Loading