From 70e1e055b962c0ebc523358f6a438cfaa1a013be Mon Sep 17 00:00:00 2001 From: swati354 Date: Mon, 3 Feb 2025 17:56:08 +0530 Subject: [PATCH 1/2] Add salesforce connector Single query for web and salesforce --- backend/danswer/configs/constants.py | 1 + .../connectors/salesforce/connector.py | 314 ++++++++---------- .../danswer/connectors/salesforce/utils.py | 28 ++ backend/danswer/document_index/vespa/index.py | 10 +- .../app/admin/connectors/salesforce/page.tsx | 67 ++-- web/src/lib/sources.ts | 5 + web/src/lib/types.ts | 4 +- 7 files changed, 216 insertions(+), 213 deletions(-) diff --git a/backend/danswer/configs/constants.py b/backend/danswer/configs/constants.py index b29d3558b84..cb0f3e955ff 100644 --- a/backend/danswer/configs/constants.py +++ b/backend/danswer/configs/constants.py @@ -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" diff --git a/backend/danswer/connectors/salesforce/connector.py b/backend/danswer/connectors/salesforce/connector.py index 03326df4efd..15ce6764e6e 100644 --- a/backend/danswer/connectors/salesforce/connector.py +++ b/backend/danswer/connectors/salesforce/connector.py @@ -1,17 +1,15 @@ import os -from collections.abc import Iterator +import json +import requests + from datetime import datetime from datetime import timezone -from typing import Any - -from simple_salesforce import Salesforce -from simple_salesforce import SFType +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 IdConnector from danswer.connectors.interfaces import LoadConnector from danswer.connectors.interfaces import PollConnector from danswer.connectors.interfaces import SecondsSinceUnixEpoch @@ -21,70 +19,70 @@ from danswer.connectors.models import Section from danswer.connectors.salesforce.utils import extract_dict_text from danswer.utils.logger import setup_logger +from danswer.connectors.salesforce.utils import clean_html -DEFAULT_PARENT_OBJECT_TYPES = ["Account"] -MAX_QUERY_LENGTH = 10000 # max query length is 20,000 characters ID_PREFIX = "SALESFORCE_" +AUTH_URL = "https://test.salesforce.com/services/oauth2/token" logger = setup_logger() -class SalesforceConnector(LoadConnector, PollConnector, IdConnector): +class SalesforceConnector(LoadConnector, PollConnector): def __init__( self, batch_size: int = INDEX_BATCH_SIZE, requested_objects: list[str] = [], ) -> None: self.batch_size = batch_size - self.sf_client: Salesforce | None = None - self.parent_object_list = ( - [obj.capitalize() for obj in requested_objects] + self.product_component_list = ( + [obj.strip() for obj in requested_objects[0].split(",")] if requested_objects - else DEFAULT_PARENT_OBJECT_TYPES + else None ) def load_credentials(self, credentials: dict[str, Any]) -> dict[str, Any] | None: - self.sf_client = Salesforce( - username=credentials["sf_username"], - password=credentials["sf_password"], - security_token=credentials["sf_security_token"], - ) + self.client_id = credentials['sf_client_id'] + self.client_secret = credentials['sf_client_secret'] + self.username = credentials['sf_username'] + self.password = credentials['sf_password'] - return None + self.access_token, self.instance_url = self._get_access_token() - def _get_sf_type_object_json(self, type_name: str) -> Any: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - sf_object = SFType( - type_name, self.sf_client.session_id, self.sf_client.sf_instance - ) - return sf_object.describe() + self.headers = { + "Authorization": f"Bearer {self.access_token}", + "Content-Type": "application/json", + } - def _get_name_from_id(self, id: str) -> str: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - try: - user_object_info = self.sf_client.query( - f"SELECT Name FROM User WHERE Id = '{id}'" - ) - name = user_object_info.get("Records", [{}])[0].get("Name", "Null User") - return name - except Exception: - logger.warning(f"Couldnt find name for object id: {id}") - return "Null User" + 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: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") salesforce_id = object_dict["Id"] danswer_salesforce_id = f"{ID_PREFIX}{salesforce_id}" - extracted_link = f"https://{self.sf_client.sf_instance}/{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("Name", "Unknown Object") + extracted_semantic_identifier = object_dict.get("Title", "Unknown Object") extracted_primary_owners = [ BasicExpertInfo( display_name=self._get_name_from_id(object_dict["LastModifiedById"]) @@ -94,7 +92,7 @@ def _convert_object_instance_to_document( doc = Document( id=danswer_salesforce_id, sections=[Section(link=extracted_link, text=extracted_object_text)], - source=DocumentSource.SALESFORCE, + source=DocumentSource.SFKBARTICLES, semantic_identifier=extracted_semantic_identifier, doc_updated_at=extracted_doc_updated_at, primary_owners=extracted_primary_owners, @@ -102,133 +100,128 @@ def _convert_object_instance_to_document( ) return doc - def _is_valid_child_object(self, child_relationship: dict) -> bool: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - - if not child_relationship["childSObject"]: - return False - if not child_relationship["relationshipName"]: - return False - - sf_type = child_relationship["childSObject"] - object_description = self._get_sf_type_object_json(sf_type) - if not object_description["queryable"]: - return False - - try: - query = f"SELECT Count() FROM {sf_type} LIMIT 1" - result = self.sf_client.query(query) - if result["totalSize"] == 0: - return False - except Exception as e: - logger.warning(f"Object type {sf_type} doesn't support query: {e}") - return False - - if child_relationship["field"]: - if child_relationship["field"] == "RelatedToId": - return False - else: - return False - - return True - - def _get_all_children_of_sf_type(self, sf_type: str) -> list[dict]: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - - object_description = self._get_sf_type_object_json(sf_type) - - children_objects: list[dict] = [] - for child_relationship in object_description["childRelationships"]: - if self._is_valid_child_object(child_relationship): - children_objects.append( - { - "relationship_name": child_relationship["relationshipName"], - "object_type": child_relationship["childSObject"], - } - ) - return children_objects - - def _get_all_fields_for_sf_type(self, sf_type: str) -> list[str]: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - - object_description = self._get_sf_type_object_json(sf_type) - - fields = [ - field.get("name") - for field in object_description["fields"] - if field.get("type", "base64") != "base64" - ] - - return fields - - def _generate_query_per_parent_type(self, parent_sf_type: str) -> Iterator[str]: + def _get_name_from_id(self, id: str) -> str: """ - This function takes in an object_type and generates query(s) designed to grab - information associated to objects of that type. - It does that by getting all the fields of the parent object type. - Then it gets all the child objects of that object type and all the fields of - those children as well. + Fetches the name of a Salesforce user based on their ID. """ - parent_fields = self._get_all_fields_for_sf_type(parent_sf_type) - child_sf_types = self._get_all_children_of_sf_type(parent_sf_type) + query = f"SELECT Name FROM User WHERE Id = '{id}'" + url = f"{self.instance_url}/services/data/v56.0/query" + params = {"q": query} - query = f"SELECT {', '.join(parent_fields)}" - for child_object_dict in child_sf_types: - fields = self._get_all_fields_for_sf_type(child_object_dict["object_type"]) - query_addition = f", \n(SELECT {', '.join(fields)} FROM {child_object_dict['relationship_name']})" + 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" - if len(query_addition) + len(query) > MAX_QUERY_LENGTH: - query += f"\n FROM {parent_sf_type}" - yield query - query = "SELECT Id" + query_addition - else: - query += query_addition + data = response.json() + records = data.get("records", []) + if not records: + logger.warning(f"No name found for ID {id}") + return "Unknown" - query += f"\n FROM {parent_sf_type}" + return records[0].get("Name", "Unknown") - yield query + 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: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - + + query = self.build_salesforce_query(self.product_component_list, start, end) doc_batch: list[Document] = [] - for parent_object_type in self.parent_object_list: - logger.debug(f"Processing: {parent_object_type}") + query_results: dict = {} - query_results: dict = {} - for query in self._generate_query_per_parent_type(parent_object_type): - if start is not None and end is not None: - if start and start.tzinfo is None: - start = start.replace(tzinfo=timezone.utc) - if end and end.tzinfo is None: - end = end.replace(tzinfo=timezone.utc) - query += f" WHERE LastModifiedDate > {start.isoformat()} AND LastModifiedDate < {end.isoformat()}" + 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 + ) - query_result = self.sf_client.query_all(query) + while url: + query_result = requests.get( + url, headers=self.headers, params=params if "q" in params else None + ) + data = query_result.json() - for record_dict in query_result["records"]: + if "records" in data: + for record_dict in data["records"]: query_results.setdefault(record_dict["Id"], {}).update(record_dict) - logger.info( - f"Number of {parent_object_type} Objects processed: {len(query_results)}" - ) + 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: + 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: @@ -237,26 +230,10 @@ def load_from_state(self) -> GenerateDocumentsOutput: def poll_source( self, start: SecondsSinceUnixEpoch, end: SecondsSinceUnixEpoch ) -> GenerateDocumentsOutput: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - start_datetime = datetime.utcfromtimestamp(start) - end_datetime = datetime.utcfromtimestamp(end) + 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) - def retrieve_all_source_ids(self) -> set[str]: - if self.sf_client is None: - raise ConnectorMissingCredentialError("Salesforce") - all_retrieved_ids: set[str] = set() - for parent_object_type in self.parent_object_list: - query = f"SELECT Id FROM {parent_object_type}" - query_result = self.sf_client.query_all(query) - all_retrieved_ids.update( - f"{ID_PREFIX}{instance_dict.get('Id', '')}" - for instance_dict in query_result["records"] - ) - - return all_retrieved_ids - if __name__ == "__main__": connector = SalesforceConnector( @@ -265,9 +242,10 @@ def retrieve_all_source_ids(self) -> set[str]: 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"], - "sf_security_token": os.environ["SF_SECURITY_TOKEN"], } ) document_batches = connector.load_from_state() diff --git a/backend/danswer/connectors/salesforce/utils.py b/backend/danswer/connectors/salesforce/utils.py index db8edf073f1..e1da77867f3 100644 --- a/backend/danswer/connectors/salesforce/utils.py +++ b/backend/danswer/connectors/salesforce/utils.py @@ -1,5 +1,7 @@ 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$" @@ -62,5 +64,31 @@ 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) + + 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 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 diff --git a/backend/danswer/document_index/vespa/index.py b/backend/danswer/document_index/vespa/index.py index a7892e40f45..e5e968fb48d 100644 --- a/backend/danswer/document_index/vespa/index.py +++ b/backend/danswer/document_index/vespa/index.py @@ -699,17 +699,17 @@ def _query_vespa(query_params: Mapping[str, str | int | float]) -> list[Inferenc if LOG_VESPA_TIMING_INFORMATION else {}, ) - + #All records including web params["hits"] = 50 filtered_hits_all = query_vespa_helper(params) - #Only Web Records + #Only Web Records and Salesforce KB articles params["hits"] = 10 - params["yql"] = params["yql"] + ' and source_type contains "web"' - filtered_hits_web = query_vespa_helper(params) + params["yql"] = params["yql"] + ' and (source_type contains "web" or source_type contains "sfkabarticles")' + filtered_hits_web_sf = query_vespa_helper(params) - filtered_hits_final = filtered_hits_web + filtered_hits_all + filtered_hits_final = filtered_hits_web_sf + filtered_hits_all inference_chunks = [_vespa_hit_to_inference_chunk(hit) for hit in filtered_hits_final] #inplace sorting based on score diff --git a/web/src/app/admin/connectors/salesforce/page.tsx b/web/src/app/admin/connectors/salesforce/page.tsx index e699f514ff6..a3d3752f100 100644 --- a/web/src/app/admin/connectors/salesforce/page.tsx +++ b/web/src/app/admin/connectors/salesforce/page.tsx @@ -102,22 +102,22 @@ const MainSection = () => { ) : ( <> - As a first step, please provide the Salesforce admin account's - username, password, and Salesforce security token. You can follow - the guide{" "} - - here - {" "} - to create get your Salesforce Security Token. + As a first step, please provide the Salesforce account's + client_id, client_secret, username and password. formBody={ <> + + { label="Salesforce Password:" type="password" /> - } validationSchema={Yup.object().shape({ + sf_client_id: Yup.string().required( + "Please enter your Salesforce Client Id" + ), + sf_client_secret: Yup.string().required( + "Please enter your Salesforce Client Secret" + ), sf_username: Yup.string().required( "Please enter your Salesforce username" ), sf_password: Yup.string().required( "Please enter your Salesforce password" ), - sf_security_token: Yup.string().required( - "Please enter your Salesforce security token" - ), })} initialValues={{ + sf_client_id: "", + sf_client_secret: "", sf_username: "", sf_password: "", - sf_security_token: "", }} onSubmit={(isSuccess) => { if (isSuccess) { @@ -175,7 +174,7 @@ const MainSection = () => { connectorIndexingStatuses={SalesforceConnectorIndexingStatuses} liveCredential={SalesforceCredential} getCredential={(credential) => - credential.credential_json.sf_security_token + credential.credential_json.sf_client_secret } onUpdate={() => mutate("/api/manage/admin/connector/indexing-status") @@ -221,30 +220,20 @@ const MainSection = () => { // formBody={<>} formBodyBuilder={TextArrayFieldBuilder({ name: "requested_objects", - label: "Specify Salesforce objects to organize by:", + label: "Specify the Product Components", subtext: ( <>
- Specify the Salesforce object types you want us to index.{" "} + Specify the product components for which you want to fetch the Salesforce Knowledge Base articles. +

+ Example: Orchestrator, Activities, Studio, Robot, Automation Hub.
- Click - - {" "} - here{" "} - - for an example of how Darwin uses the objects.

- If unsure, don't specify any objects and Darwin will - default to indexing by 'Account'. + By default, it will fetch articles for all the product components.

- Hint: Use the singular form of the object name (e.g., - 'Opportunity' instead of 'Opportunities'). + Hint: Use the exact product component name for accurate results. ), })} @@ -267,8 +256,8 @@ const MainSection = () => { ) : ( Please provide all Salesforce info in Step 1 first! Once you're - done with that, you can then specify which Salesforce objects you want - to make searchable. + done with that, you can then specify the product components for which + you want to fetch the Salesforce Knowledge Base articles. )} diff --git a/web/src/lib/sources.ts b/web/src/lib/sources.ts index f141e42edb5..86ca83e0938 100644 --- a/web/src/lib/sources.ts +++ b/web/src/lib/sources.ts @@ -172,6 +172,11 @@ const SOURCE_METADATA_MAP: SourceMap = { displayName: "Salesforce", category: SourceCategory.AppConnection, }, + sfkbarticles: { + icon: SalesforceIcon, + displayName: "sfkbarticles", + category: SourceCategory.AppConnection, + }, sharepoint: { icon: SharepointIcon, displayName: "Sharepoint", diff --git a/web/src/lib/types.ts b/web/src/lib/types.ts index fe252afbf28..e563f0a830b 100644 --- a/web/src/lib/types.ts +++ b/web/src/lib/types.ts @@ -51,6 +51,7 @@ export type ValidSources = | "loopio" | "dropbox" | "salesforce" + | "sfkbarticles" | "sharepoint" | "teams" | "zendesk" @@ -452,9 +453,10 @@ export interface OCICredentialJson { secret_access_key: string; } export interface SalesforceCredentialJson { + sf_client_id: string; + sf_client_secret: string; sf_username: string; sf_password: string; - sf_security_token: string; } export interface SharepointCredentialJson { From bbe26a7cecafa2517270a211c328d4ea81326467 Mon Sep 17 00:00:00 2001 From: swati354 Date: Mon, 17 Feb 2025 00:20:59 +0530 Subject: [PATCH 2/2] Add sfkbarticles connector Update code comments --- backend/danswer/connectors/factory.py | 2 + .../connectors/salesforce/connector.py | 316 ++++++++++-------- .../danswer/connectors/salesforce/utils.py | 30 +- .../connectors/sfkbarticles/__init__.py | 0 .../connectors/sfkbarticles/connector.py | 253 ++++++++++++++ .../danswer/connectors/sfkbarticles/utils.py | 94 ++++++ backend/danswer/prompts/prompt_utils.py | 7 +- backend/danswer/utils/text_processing.py | 11 +- .../app/admin/connectors/salesforce/page.tsx | 67 ++-- .../admin/connectors/sfkbarticles/page.tsx | 279 ++++++++++++++++ web/src/lib/sources.ts | 2 +- web/src/lib/types.ts | 11 + 12 files changed, 863 insertions(+), 209 deletions(-) create mode 100644 backend/danswer/connectors/sfkbarticles/__init__.py create mode 100644 backend/danswer/connectors/sfkbarticles/connector.py create mode 100644 backend/danswer/connectors/sfkbarticles/utils.py create mode 100644 web/src/app/admin/connectors/sfkbarticles/page.tsx diff --git a/backend/danswer/connectors/factory.py b/backend/danswer/connectors/factory.py index 1a3d605d3a5..2f885bf2b7d 100644 --- a/backend/danswer/connectors/factory.py +++ b/backend/danswer/connectors/factory.py @@ -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 @@ -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, diff --git a/backend/danswer/connectors/salesforce/connector.py b/backend/danswer/connectors/salesforce/connector.py index 15ce6764e6e..436b016dbe7 100644 --- a/backend/danswer/connectors/salesforce/connector.py +++ b/backend/danswer/connectors/salesforce/connector.py @@ -1,15 +1,17 @@ import os -import json -import requests - +from collections.abc import Iterator from datetime import datetime from datetime import timezone -from typing import Any, Tuple +from typing import Any + +from simple_salesforce import Salesforce +from simple_salesforce import SFType 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 IdConnector from danswer.connectors.interfaces import LoadConnector from danswer.connectors.interfaces import PollConnector from danswer.connectors.interfaces import SecondsSinceUnixEpoch @@ -19,70 +21,70 @@ from danswer.connectors.models import Section from danswer.connectors.salesforce.utils import extract_dict_text from danswer.utils.logger import setup_logger -from danswer.connectors.salesforce.utils import clean_html +DEFAULT_PARENT_OBJECT_TYPES = ["Account"] +MAX_QUERY_LENGTH = 10000 # max query length is 20,000 characters ID_PREFIX = "SALESFORCE_" -AUTH_URL = "https://test.salesforce.com/services/oauth2/token" logger = setup_logger() -class SalesforceConnector(LoadConnector, PollConnector): +class SalesforceConnector(LoadConnector, PollConnector, IdConnector): 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(",")] + self.sf_client: Salesforce | None = None + self.parent_object_list = ( + [obj.capitalize() for obj in requested_objects] if requested_objects - else None + else DEFAULT_PARENT_OBJECT_TYPES ) 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.sf_client = Salesforce( + username=credentials["sf_username"], + password=credentials["sf_password"], + security_token=credentials["sf_security_token"], + ) - self.headers = { - "Authorization": f"Bearer {self.access_token}", - "Content-Type": "application/json", - } + return None - 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.") + def _get_sf_type_object_json(self, type_name: str) -> Any: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + sf_object = SFType( + type_name, self.sf_client.session_id, self.sf_client.sf_instance + ) + return sf_object.describe() - data = response.json() - logger.info("Successfully authenticated with Salesforce.") - return data["access_token"], data["instance_url"] + def _get_name_from_id(self, id: str) -> str: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + try: + user_object_info = self.sf_client.query( + f"SELECT Name FROM User WHERE Id = '{id}'" + ) + name = user_object_info.get("Records", [{}])[0].get("Name", "Null User") + return name + except Exception: + logger.warning(f"Couldnt find name for object id: {id}") + return "Null User" def _convert_object_instance_to_document( self, object_dict: dict[str, Any] ) -> Document: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") salesforce_id = object_dict["Id"] danswer_salesforce_id = f"{ID_PREFIX}{salesforce_id}" - extracted_link = f"{self.instance_url}/{salesforce_id}" + extracted_link = f"https://{self.sf_client.sf_instance}/{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_semantic_identifier = object_dict.get("Name", "Unknown Object") extracted_primary_owners = [ BasicExpertInfo( display_name=self._get_name_from_id(object_dict["LastModifiedById"]) @@ -92,7 +94,7 @@ def _convert_object_instance_to_document( doc = Document( id=danswer_salesforce_id, sections=[Section(link=extracted_link, text=extracted_object_text)], - source=DocumentSource.SFKBARTICLES, + source=DocumentSource.SALESFORCE, semantic_identifier=extracted_semantic_identifier, doc_updated_at=extracted_doc_updated_at, primary_owners=extracted_primary_owners, @@ -100,128 +102,133 @@ def _convert_object_instance_to_document( ) 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} + def _is_valid_child_object(self, child_relationship: dict) -> bool: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + + if not child_relationship["childSObject"]: + return False + if not child_relationship["relationshipName"]: + return False + + sf_type = child_relationship["childSObject"] + object_description = self._get_sf_type_object_json(sf_type) + if not object_description["queryable"]: + return False + + try: + query = f"SELECT Count() FROM {sf_type} LIMIT 1" + result = self.sf_client.query(query) + if result["totalSize"] == 0: + return False + except Exception as e: + logger.warning(f"Object type {sf_type} doesn't support query: {e}") + return False + + if child_relationship["field"]: + if child_relationship["field"] == "RelatedToId": + return False + else: + return False - 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" + return True - data = response.json() - records = data.get("records", []) - if not records: - logger.warning(f"No name found for ID {id}") - return "Unknown" + def _get_all_children_of_sf_type(self, sf_type: str) -> list[dict]: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") - return records[0].get("Name", "Unknown") + object_description = self._get_sf_type_object_json(sf_type) - def build_salesforce_query(self, - parent_objects: list[str], - start: datetime | None = None, - end: datetime | None = None) -> str: + children_objects: list[dict] = [] + for child_relationship in object_description["childRelationships"]: + if self._is_valid_child_object(child_relationship): + children_objects.append( + { + "relationship_name": child_relationship["relationshipName"], + "object_type": child_relationship["childSObject"], + } + ) + return children_objects + + def _get_all_fields_for_sf_type(self, sf_type: str) -> list[str]: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + + object_description = self._get_sf_type_object_json(sf_type) + + fields = [ + field.get("name") + for field in object_description["fields"] + if field.get("type", "base64") != "base64" + ] + + return fields + + def _generate_query_per_parent_type(self, parent_sf_type: str) -> Iterator[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. + This function takes in an object_type and generates query(s) designed to grab + information associated to objects of that type. + It does that by getting all the fields of the parent object type. + Then it gets all the child objects of that object type and all the fields of + those children as well. """ - 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 + parent_fields = self._get_all_fields_for_sf_type(parent_sf_type) + child_sf_types = self._get_all_children_of_sf_type(parent_sf_type) + + query = f"SELECT {', '.join(parent_fields)}" + for child_object_dict in child_sf_types: + fields = self._get_all_fields_for_sf_type(child_object_dict["object_type"]) + query_addition = f", \n(SELECT {', '.join(fields)} FROM {child_object_dict['relationship_name']})" + + if len(query_addition) + len(query) > MAX_QUERY_LENGTH: + query += f"\n FROM {parent_sf_type}" + yield query + query = "SELECT Id" + query_addition + else: + query += query_addition + + query += f"\n FROM {parent_sf_type}" + + yield 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) + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + doc_batch: list[Document] = [] - query_results: dict = {} + for parent_object_type in self.parent_object_list: + logger.debug(f"Processing: {parent_object_type}") - 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 - ) + query_results: dict = {} + for query in self._generate_query_per_parent_type(parent_object_type): + if start is not None and end is not None: + if start and start.tzinfo is None: + start = start.replace(tzinfo=timezone.utc) + if end and end.tzinfo is None: + end = end.replace(tzinfo=timezone.utc) + query += f" WHERE LastModifiedDate > {start.isoformat()} AND LastModifiedDate < {end.isoformat()}" - while url: - query_result = requests.get( - url, headers=self.headers, params=params if "q" in params else None - ) - data = query_result.json() + query_result = self.sf_client.query_all(query) - if "records" in data: - for record_dict in data["records"]: + for record_dict in query_result["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) + logger.info( + f"Number of {parent_object_type} Objects processed: {len(query_results)}" ) - if len(doc_batch) > self.batch_size: + + 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: @@ -230,10 +237,26 @@ def load_from_state(self) -> GenerateDocumentsOutput: 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) + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + start_datetime = datetime.utcfromtimestamp(start) + end_datetime = datetime.utcfromtimestamp(end) return self._fetch_from_salesforce(start=start_datetime, end=end_datetime) + def retrieve_all_source_ids(self) -> set[str]: + if self.sf_client is None: + raise ConnectorMissingCredentialError("Salesforce") + all_retrieved_ids: set[str] = set() + for parent_object_type in self.parent_object_list: + query = f"SELECT Id FROM {parent_object_type}" + query_result = self.sf_client.query_all(query) + all_retrieved_ids.update( + f"{ID_PREFIX}{instance_dict.get('Id', '')}" + for instance_dict in query_result["records"] + ) + + return all_retrieved_ids + if __name__ == "__main__": connector = SalesforceConnector( @@ -242,11 +265,10 @@ def poll_source( 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"], + "sf_security_token": os.environ["SF_SECURITY_TOKEN"], } ) document_batches = connector.load_from_state() - print(next(document_batches)) + print(next(document_batches)) \ No newline at end of file diff --git a/backend/danswer/connectors/salesforce/utils.py b/backend/danswer/connectors/salesforce/utils.py index e1da77867f3..a9dda88ee63 100644 --- a/backend/danswer/connectors/salesforce/utils.py +++ b/backend/danswer/connectors/salesforce/utils.py @@ -1,7 +1,5 @@ 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$" @@ -64,31 +62,5 @@ 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) - - 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 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 + return natural_language_dict \ No newline at end of file diff --git a/backend/danswer/connectors/sfkbarticles/__init__.py b/backend/danswer/connectors/sfkbarticles/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/backend/danswer/connectors/sfkbarticles/connector.py b/backend/danswer/connectors/sfkbarticles/connector.py new file mode 100644 index 00000000000..7b7be3a2553 --- /dev/null +++ b/backend/danswer/connectors/sfkbarticles/connector.py @@ -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)) diff --git a/backend/danswer/connectors/sfkbarticles/utils.py b/backend/danswer/connectors/sfkbarticles/utils.py new file mode 100644 index 00000000000..a15b4b61875 --- /dev/null +++ b/backend/danswer/connectors/sfkbarticles/utils.py @@ -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 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 diff --git a/backend/danswer/prompts/prompt_utils.py b/backend/danswer/prompts/prompt_utils.py index 6d7bddeec95..117de69cf9a 100644 --- a/backend/danswer/prompts/prompt_utils.py +++ b/backend/danswer/prompts/prompt_utils.py @@ -178,9 +178,10 @@ def drop_messages_history_overflow( final_msgs = [final_msg] # Start dropping from the history if necessary - ind_prev_msg_start = find_last_index( - token_counts, max_prompt_tokens=max_allowed_tokens - ) + # ind_prev_msg_start = find_last_index( + # token_counts, max_prompt_tokens=max_allowed_tokens + # ) + ind_prev_msg_start = 0 if system_msg and ind_prev_msg_start <= len(history_msgs): final_messages.append(system_msg) diff --git a/backend/danswer/utils/text_processing.py b/backend/danswer/utils/text_processing.py index b0fbcdfa1e9..6083ffded15 100644 --- a/backend/danswer/utils/text_processing.py +++ b/backend/danswer/utils/text_processing.py @@ -20,7 +20,16 @@ def decode_escapes(s: str) -> str: def decode_match(match: re.Match) -> str: - return codecs.decode(match.group(0), "unicode-escape") + matched_str = match.group(0) + + # Only double escape non-Unicode sequences + if matched_str.startswith("\\U") and not re.match(r"\\U[0-9a-fA-F]{8}", matched_str): + return matched_str + + try: + return codecs.decode(matched_str, "unicode-escape") + except UnicodeDecodeError: + return matched_str return ESCAPE_SEQUENCE_RE.sub(decode_match, s) diff --git a/web/src/app/admin/connectors/salesforce/page.tsx b/web/src/app/admin/connectors/salesforce/page.tsx index a3d3752f100..e699f514ff6 100644 --- a/web/src/app/admin/connectors/salesforce/page.tsx +++ b/web/src/app/admin/connectors/salesforce/page.tsx @@ -102,22 +102,22 @@ const MainSection = () => { ) : ( <> - As a first step, please provide the Salesforce account's - client_id, client_secret, username and password. + As a first step, please provide the Salesforce admin account's + username, password, and Salesforce security token. You can follow + the guide{" "} + + here + {" "} + to create get your Salesforce Security Token. formBody={ <> - - { label="Salesforce Password:" type="password" /> + } validationSchema={Yup.object().shape({ - sf_client_id: Yup.string().required( - "Please enter your Salesforce Client Id" - ), - sf_client_secret: Yup.string().required( - "Please enter your Salesforce Client Secret" - ), sf_username: Yup.string().required( "Please enter your Salesforce username" ), sf_password: Yup.string().required( "Please enter your Salesforce password" ), + sf_security_token: Yup.string().required( + "Please enter your Salesforce security token" + ), })} initialValues={{ - sf_client_id: "", - sf_client_secret: "", sf_username: "", sf_password: "", + sf_security_token: "", }} onSubmit={(isSuccess) => { if (isSuccess) { @@ -174,7 +175,7 @@ const MainSection = () => { connectorIndexingStatuses={SalesforceConnectorIndexingStatuses} liveCredential={SalesforceCredential} getCredential={(credential) => - credential.credential_json.sf_client_secret + credential.credential_json.sf_security_token } onUpdate={() => mutate("/api/manage/admin/connector/indexing-status") @@ -220,20 +221,30 @@ const MainSection = () => { // formBody={<>} formBodyBuilder={TextArrayFieldBuilder({ name: "requested_objects", - label: "Specify the Product Components", + label: "Specify Salesforce objects to organize by:", subtext: ( <>
- Specify the product components for which you want to fetch the Salesforce Knowledge Base articles. -
+ Specify the Salesforce object types you want us to index.{" "}
- Example: Orchestrator, Activities, Studio, Robot, Automation Hub.
+ Click + + {" "} + here{" "} + + for an example of how Darwin uses the objects.

- By default, it will fetch articles for all the product components. + If unsure, don't specify any objects and Darwin will + default to indexing by 'Account'.

- Hint: Use the exact product component name for accurate results. + Hint: Use the singular form of the object name (e.g., + 'Opportunity' instead of 'Opportunities'). ), })} @@ -256,8 +267,8 @@ const MainSection = () => { ) : ( Please provide all Salesforce info in Step 1 first! Once you're - done with that, you can then specify the product components for which - you want to fetch the Salesforce Knowledge Base articles. + done with that, you can then specify which Salesforce objects you want + to make searchable. )} diff --git a/web/src/app/admin/connectors/sfkbarticles/page.tsx b/web/src/app/admin/connectors/sfkbarticles/page.tsx new file mode 100644 index 00000000000..bd07892e537 --- /dev/null +++ b/web/src/app/admin/connectors/sfkbarticles/page.tsx @@ -0,0 +1,279 @@ +"use client"; + +import * as Yup from "yup"; +import { TrashIcon, SalesforceIcon } from "@/components/icons/icons"; // Make sure you have a Document360 icon +import { errorHandlingFetcher as fetcher } from "@/lib/fetcher"; +import useSWR, { useSWRConfig } from "swr"; +import { LoadingAnimation } from "@/components/Loading"; +import { HealthCheckBanner } from "@/components/health/healthcheck"; +import { + SfKbArticlesConfig, + SfKbArticlesCredentialJson, + ConnectorIndexingStatus, + Credential, +} from "@/lib/types"; // Modify or create these types as required +import { adminDeleteCredential, linkCredential } from "@/lib/credential"; +import { CredentialForm } from "@/components/admin/connectors/CredentialForm"; +import { + TextFormField, + TextArrayFieldBuilder, +} from "@/components/admin/connectors/Field"; +import { ConnectorsTable } from "@/components/admin/connectors/table/ConnectorsTable"; +import { ConnectorForm } from "@/components/admin/connectors/ConnectorForm"; +import { usePublicCredentials } from "@/lib/hooks"; +import { AdminPageTitle } from "@/components/admin/Title"; +import { Card, Text, Title } from "@tremor/react"; + +const MainSection = () => { + const { mutate } = useSWRConfig(); + const { + data: connectorIndexingStatuses, + isLoading: isConnectorIndexingStatusesLoading, + error: isConnectorIndexingStatusesError, + } = useSWR[]>( + "/api/manage/admin/connector/indexing-status", + fetcher + ); + + const { + data: credentialsData, + isLoading: isCredentialsLoading, + error: isCredentialsError, + refreshCredentials, + } = usePublicCredentials(); + + if ( + (!connectorIndexingStatuses && isConnectorIndexingStatusesLoading) || + (!credentialsData && isCredentialsLoading) + ) { + return ; + } + + if (isConnectorIndexingStatusesError || !connectorIndexingStatuses) { + return
Failed to load connectors
; + } + + if (isCredentialsError || !credentialsData) { + return
Failed to load credentials
; + } + + const SalesforceConnectorIndexingStatuses: ConnectorIndexingStatus< + SfKbArticlesConfig, + SfKbArticlesCredentialJson + >[] = connectorIndexingStatuses.filter( + (connectorIndexingStatus) => + connectorIndexingStatus.connector.source === "salesforce" + ); + + const SfKbArticlesCredential: Credential | undefined = + credentialsData.find( + (credential) => credential.credential_json?.sf_username + ); + + return ( + <> + + The Salesforce Knowledge Base Articles connector allows you to index and search through your + Salesforce Knowledge Base. Once setup, all indicated Salesforce data will + be queryable within Darwin. + + + + Step 1: Provide Salesforce credentials + + {SfKbArticlesCredential ? ( + <> +
+ Existing SalesForce Username: + + {SfKbArticlesCredential.credential_json.sf_username} + + +
+ + ) : ( + <> + + As a first step, please provide the Salesforce account's + client_id, client_secret, username and password. + + + + formBody={ + <> + + + + + + } + validationSchema={Yup.object().shape({ + sf_client_id: Yup.string().required( + "Please enter your Salesforce Client Id" + ), + sf_client_secret: Yup.string().required( + "Please enter your Salesforce Client Secret" + ), + sf_username: Yup.string().required( + "Please enter your Salesforce username" + ), + sf_password: Yup.string().required( + "Please enter your Salesforce password" + ), + })} + initialValues={{ + sf_client_id: "", + sf_client_secret: "", + sf_username: "", + sf_password: "", + }} + onSubmit={(isSuccess) => { + if (isSuccess) { + refreshCredentials(); + } + }} + /> + + + )} + + + Step 2: Manage Salesforce KB Articles Connector + + + {SalesforceConnectorIndexingStatuses.length > 0 && ( + <> + + The latest state of your Salesforce objects are fetched every 10 + minutes. + +
+ + connectorIndexingStatuses={SalesforceConnectorIndexingStatuses} + liveCredential={SfKbArticlesCredential} + getCredential={(credential) => + credential.credential_json.sf_client_secret + } + onUpdate={() => + mutate("/api/manage/admin/connector/indexing-status") + } + onCredentialLink={async (connectorId) => { + if (SfKbArticlesCredential) { + await linkCredential(connectorId, SfKbArticlesCredential.id); + mutate("/api/manage/admin/connector/indexing-status"); + } + }} + specialColumns={[ + { + header: "Connectors", + key: "connectors", + getValue: (ccPairStatus) => { + const connectorConfig = + ccPairStatus.connector.connector_specific_config; + return `${connectorConfig.requested_objects}`; + }, + }, + ]} + includeName + /> +
+ + )} + + {SfKbArticlesCredential ? ( + + + nameBuilder={(values) => + values.requested_objects && values.requested_objects.length > 0 + ? `SfKbArticles-${values.requested_objects.join("-")}` + : "SfKbArticles" + } + ccPairNameBuilder={(values) => + values.requested_objects && values.requested_objects.length > 0 + ? `SfKbArticles-${values.requested_objects.join("-")}` + : "SfKbArticles" + } + source="sfkbarticles" + inputType="poll" + // formBody={<>} + formBodyBuilder={TextArrayFieldBuilder({ + name: "requested_objects", + label: "Specify the Product Components", + subtext: ( + <> +
+ Specify the product components for which you want to fetch the Salesforce Knowledge Base articles. +
+
+ Example: Orchestrator, Activities, Studio, Robot, Automation Hub. +
+
+ By default, it will fetch articles for all the product components. +
+
+ Hint: Use the exact product component name for accurate results. + + ), + })} + validationSchema={Yup.object().shape({ + requested_objects: Yup.array() + .of( + Yup.string().required( + "Salesforce Product Component names must be strings" + ) + ) + .required(), + })} + initialValues={{ + requested_objects: [], + }} + credentialId={SfKbArticlesCredential.id} + refreshFreq={10 * 60} // 10 minutes + /> +
+ ) : ( + + Please provide all Salesforce info in Step 1 first! Once you're + done with that, you can then specify the product components for which + you want to fetch the Salesforce Knowledge Base articles. + + )} + + ); +}; + +export default function Page() { + return ( +
+
+ +
+ + } title="Salesforce KB Articles" /> + + +
+ ); +} diff --git a/web/src/lib/sources.ts b/web/src/lib/sources.ts index 86ca83e0938..6c93d412a69 100644 --- a/web/src/lib/sources.ts +++ b/web/src/lib/sources.ts @@ -174,7 +174,7 @@ const SOURCE_METADATA_MAP: SourceMap = { }, sfkbarticles: { icon: SalesforceIcon, - displayName: "sfkbarticles", + displayName: "SfKbArticles", category: SourceCategory.AppConnection, }, sharepoint: { diff --git a/web/src/lib/types.ts b/web/src/lib/types.ts index e563f0a830b..ea6f3c7defc 100644 --- a/web/src/lib/types.ts +++ b/web/src/lib/types.ts @@ -145,6 +145,10 @@ export interface SalesforceConfig { requested_objects?: string[]; } +export interface SfKbArticlesConfig { + requested_objects?: string[]; +} + export interface SharepointConfig { sites?: string[]; } @@ -452,7 +456,14 @@ export interface OCICredentialJson { access_key_id: string; secret_access_key: string; } + export interface SalesforceCredentialJson { + sf_username: string; + sf_password: string; + sf_security_token: string; +} + +export interface SfKbArticlesCredentialJson { sf_client_id: string; sf_client_secret: string; sf_username: string;