From 592922a30e76ae52075b257785b79d3acd566cdc Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Thu, 7 Dec 2023 09:57:22 +0800 Subject: [PATCH 01/10] Converted smart search from pinecone to postgre --- .env.example | 12 ++++ requirements.txt | 4 +- src/semantic_search/semantic_search/config.py | 15 +++++ .../external_services/postgre_vector.py | 38 ++++++++++++ .../external_services/slack_api.py | 5 +- .../handle_indexation_tasks.py | 3 +- .../semantic_search/load_messages.py | 59 +++++++++---------- src/semantic_search/semantic_search/query.py | 50 +++++++++------- src/services/api_service.py | 2 + 9 files changed, 132 insertions(+), 56 deletions(-) create mode 100644 src/semantic_search/semantic_search/external_services/postgre_vector.py diff --git a/.env.example b/.env.example index 0b2dfa9..3625b30 100644 --- a/.env.example +++ b/.env.example @@ -17,3 +17,15 @@ LOG_LEVEL=DEBUG STANDALONE=true SLACK_USER_ID=U01JZQZQZQZ USE_FALLBACK=false + +# PINECONE +PINECONE_INDEX=index-name +PINECONE_ENVIRONMENT=environment +PINECONE_KEY=key + +# POSTGRE +POSTGRE_HOST=localhost +POSTGRE_PORT=5432 +POSTGRE_DATABASE=db-name +POSTGRE_USER=postgres +POSTGRE_PASSWORD=super-secret \ No newline at end of file diff --git a/requirements.txt b/requirements.txt index 0ea6bf1..c3628e7 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,4 +3,6 @@ python-dotenv==1.0.0 slack-bolt==1.18.0 openai==0.28 gunicorn==20.1.0 -replicate==0.18.1 \ No newline at end of file +replicate==0.18.1 +psycopg[binary] +psycopg2-binary \ No newline at end of file diff --git a/src/semantic_search/semantic_search/config.py b/src/semantic_search/semantic_search/config.py index 53e6345..9c077bf 100644 --- a/src/semantic_search/semantic_search/config.py +++ b/src/semantic_search/semantic_search/config.py @@ -54,3 +54,18 @@ def get_api_shared_secret() -> str: def is_standalone() -> bool: return os.environ.get('STANDALONE') == 'true' + +def get_postgre_host()-> str : + return os.environ.get('POSTGRE_HOST') + +def get_postgre_port()-> str : + return os.environ.get('POSTGRE_PORT') + +def get_postgre_database()-> str : + return os.environ.get('POSTGRE_DATABASE') + +def get_postgre_user()-> str : + return os.environ.get('POSTGRE_USER') + +def get_postgre_password()-> str : + return os.environ.get('POSTGRE_PASSWORD') diff --git a/src/semantic_search/semantic_search/external_services/postgre_vector.py b/src/semantic_search/semantic_search/external_services/postgre_vector.py new file mode 100644 index 0000000..1280cdb --- /dev/null +++ b/src/semantic_search/semantic_search/external_services/postgre_vector.py @@ -0,0 +1,38 @@ +from pgvector.psycopg2 import register_vector +import psycopg2 + +from ..config import get_postgre_host, get_postgre_port, get_postgre_database, get_postgre_user, get_postgre_password + +conn = None +try: + conn = psycopg2.connect( + host=get_postgre_host(), + database=get_postgre_database(), + user=get_postgre_user(), + password=get_postgre_password(), + port=get_postgre_port()) + + cur = conn.cursor() + + cur.execute('CREATE EXTENSION IF NOT EXISTS vector') + register_vector(cur) + + # cur.execute('DROP TABLE IF EXISTS embedding') + cur.execute('CREATE TABLE IF NOT EXISTS embedding (id bigserial PRIMARY KEY, namespace text, chunk_id text, metadata text, values vector)') + + conn.commit() +except (Exception, psycopg2.DatabaseError) as error: + print(error) + + +def get_postgre_cursor(): + return cur + + +def postgre_commit(): + if conn is not None: + conn.commit() + +def postgre_excute(postgre_cursor): + if conn is not None: + postgre_cursor.excute(postgre_cursor) \ No newline at end of file diff --git a/src/semantic_search/semantic_search/external_services/slack_api.py b/src/semantic_search/semantic_search/external_services/slack_api.py index 5e1ce5b..294cefe 100644 --- a/src/semantic_search/semantic_search/external_services/slack_api.py +++ b/src/semantic_search/semantic_search/external_services/slack_api.py @@ -114,8 +114,9 @@ def slack_names_map(team_id): def load_previous_messages(team_id: str, channel_id: str, last_message_id: str, number: int): - (messages, ) = load_previous_messages_with_pointer(team_id, channel_id, last_message_id, number) - return messages[-number:] + result = load_previous_messages_with_pointer(team_id, channel_id, last_message_id, number) + messages = result[0][-number:] + return messages def load_subsequent_messages(team_id: str, channel_id: str, first_message_id: str, number: int): diff --git a/src/semantic_search/semantic_search/handle_indexation_tasks.py b/src/semantic_search/semantic_search/handle_indexation_tasks.py index a927903..ff0bfb1 100644 --- a/src/semantic_search/semantic_search/handle_indexation_tasks.py +++ b/src/semantic_search/semantic_search/handle_indexation_tasks.py @@ -4,6 +4,7 @@ from flask import Flask, request, jsonify from .config import get_google_tasks_service_account from .external_services.pinecone import get_pinecone_index +from .external_services.postgre_vector import get_postgre_cursor from .google_tasks import queue_task from .load_messages import index_messages from .external_services.slack_api import load_previous_messages_with_pointer @@ -43,7 +44,7 @@ def handle_task(): [messages, next_last_message, start_from] = load_previous_messages_with_pointer(namespace, channel_id, last_message_id, BULK_SIZE) logging.info(f"Task: {task_id}, Iteration Number: {iteration_number}") logging.info(f"Task: {task_id}, Number of Actual Messages: {len(messages)}") - index_messages(channel_id, messages, start_from, get_pinecone_index(), namespace) + index_messages(channel_id, messages, start_from, get_postgre_cursor(), namespace) if next_last_message is not None: queue_task({ 'task_id': task_id, diff --git a/src/semantic_search/semantic_search/load_messages.py b/src/semantic_search/semantic_search/load_messages.py index 10bba2c..465e96d 100644 --- a/src/semantic_search/semantic_search/load_messages.py +++ b/src/semantic_search/semantic_search/load_messages.py @@ -4,12 +4,15 @@ from .config import CONTEXT_LENGTH from .external_services.pinecone import get_pinecone_index +from .external_services.postgre_vector import get_postgre_cursor, postgre_commit from .external_services.openai import create_embeddings, summarize_thread_with_chat_gpt_3_5 import datetime from .external_services.slack_api import fetch_thread_messages, fetch_channel_messages, is_thread, \ is_actual_message, \ slack_names_map, filter_messages, load_previous_messages, load_subsequent_messages +import numpy as np +import json class Embedding: def __init__(self, channel_id, id, text, ts, thread_ts=None, author_id=None): @@ -118,12 +121,12 @@ def attach_header(embeddings: List[Embedding], header: Embedding) -> List[Embedd return part_with_header + [embedding.add_header(header) for embedding in part_without_header] -def index_messages(channel_id, messages, start_from, pinecone_index, pinecone_namespace): +def index_messages(channel_id, messages, start_from, postgre_cursor, namespace): total_messages = len(messages) logging.info("Replacing User IDs with User Names in the messages") embeddings = generate_embeddings(channel_id, messages) - embeddings = replace_ids_with_names(embeddings, team_id=pinecone_namespace) + embeddings = replace_ids_with_names(embeddings, team_id=namespace) embeddings = enrich_with_datetime(embeddings) embeddings_without_context = embeddings embeddings = enrich_with_adjacent_messages(embeddings) @@ -135,9 +138,9 @@ def index_messages(channel_id, messages, start_from, pinecone_index, pinecone_na if is_thread(message): logging.info( f"{counter + 1}/{total_messages} Appending thread messages for {message['ts']} : {message['thread_ts']}") - thread_messages = filter_messages(fetch_thread_messages(pinecone_namespace, channel_id, message["thread_ts"])) + thread_messages = filter_messages(fetch_thread_messages(namespace, channel_id, message["thread_ts"])) thread_embeddings = generate_embeddings(channel_id, thread_messages) - thread_embeddings = replace_ids_with_names(thread_embeddings, team_id=pinecone_namespace) + thread_embeddings = replace_ids_with_names(thread_embeddings, team_id=namespace) thread_embeddings = enrich_with_datetime(thread_embeddings) thread_header = thread_embeddings[0] raw_messages_for_summary = list(map(lambda e: e.text, thread_embeddings)) @@ -164,54 +167,48 @@ def index_messages(channel_id, messages, start_from, pinecone_index, pinecone_na messages_for_embedding = list(filter(lambda emb_t: len(emb_t.text) != 0, messages_for_embedding)) logging.info(f"Removed empty messages, {str(len(messages_for_embedding))} messages left") - insert_pinecone_embeddings(messages_for_embedding, pinecone_index, pinecone_namespace) + insert_db_embeddings(messages_for_embedding, postgre_cursor, namespace) -def index_whole_channel(pinecone_namespace, channel_id): +def index_whole_channel(namespace, channel_id): logging.info(f"Fetching all messages from {channel_id} channel") - messages = list(reversed(fetch_channel_messages(pinecone_namespace, channel_id))) + messages = list(reversed(fetch_channel_messages(namespace, channel_id))) logging.info(f"Loaded {str(len(messages))} messages") messages = filter_messages(messages) total_messages = len(messages) logging.info(f"Filtering out service messages, left {str(total_messages)} messages") - index_messages(channel_id, messages, 0, get_pinecone_index(), pinecone_namespace) + index_messages(channel_id, messages, 0, get_postgre_cursor(), namespace) -def insert_pinecone_embeddings(messages_for_embedding: List[Embedding], pinecone_index, pinecone_namespace): +def insert_db_embeddings(messages_for_embedding: List[Embedding], postgre_cursor, namespace): logging.info("Starting embeddings creation for the generated messages") chunk_size = 30 # for OpenAI embedding_chunks = [messages_for_embedding[i:i + chunk_size] for i in range(0, len(messages_for_embedding), chunk_size)] counter = 0 for chunk in embedding_chunks: - logging.info(f"Inserting a chunk of Pinecone embeddings: [{counter} - {counter + len(chunk) - 1}]") + logging.info(f"Inserting a chunk of Postgre embeddings: [{counter} - {counter + len(chunk) - 1}]") counter += len(chunk) try: embeddings = create_embeddings([embedding_message.text for embedding_message in chunk]) - items = [] for i in range(len(chunk)): - items.append({ - 'id': chunk[i].id, - 'values': embeddings[i], - 'metadata': chunk[i].to_metadata() - }) - - pinecone_index.upsert( - vectors=items, - namespace=pinecone_namespace - ) + metadata = json.dumps(chunk[i].to_metadata()) + postgre_cursor.execute('INSERT INTO embedding (namespace, chunk_id, metadata, values) VALUES (%s, %s, %s, %s)', (namespace, chunk[i].id, metadata, embeddings[i])) + + postgre_commit() except: logging.exception("Couldn't insert embeddings") -def delete_pinecone_embedding(embeddings: list[Embedding], pinecone_index, pinecone_namespace): +def delete_db_embedding(embeddings: List[Embedding], postgre_cursor, namespace): ids = list(map(lambda emb: emb.id, embeddings)) logging.info(f"Deleting embeddings for {str(ids)}") - pinecone_index.delete(ids=ids, namespace=pinecone_namespace) - + ids_to_delete_str = ", ".join(map(str, ids)) + postgre_cursor.execute('DELETE FROM embedding WHERE chunk_id IN (%s) AND namespace=%s',(ids_to_delete_str, namespace)) + postgre_commit() def handle_message_update_and_reindex(body): event = body['event'] @@ -224,15 +221,15 @@ def handle_message_update_and_reindex(body): if not is_actual_message(message): return embedding = slack_message_to_embedding(channel_id, message) - delete_pinecone_embedding([embedding], get_pinecone_index(), team_id) + delete_db_embedding([embedding], get_postgre_cursor(), team_id) if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_pinecone_index(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgre_cursor(), team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_pinecone_index(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgre_cursor(), team_id) return if 'subtype' in event and event['subtype'] == 'message_changed': # processing a message update @@ -242,12 +239,12 @@ def handle_message_update_and_reindex(body): return if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_pinecone_index(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgre_cursor(), team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH)[1:] # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_pinecone_index(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgre_cursor(), team_id) return if 'subtype' not in event: message = event @@ -260,9 +257,9 @@ def handle_message_update_and_reindex(body): if not is_actual_message(message): return embeddings = generate_embedding_for_message(team_id, channel_id, message_id, thread_ts) - insert_pinecone_embeddings( + insert_db_embeddings( embeddings, - get_pinecone_index(), + get_postgre_cursor(), team_id ) diff --git a/src/semantic_search/semantic_search/query.py b/src/semantic_search/semantic_search/query.py index 94220c2..017e490 100644 --- a/src/semantic_search/semantic_search/query.py +++ b/src/semantic_search/semantic_search/query.py @@ -4,7 +4,9 @@ import uuid from datetime import date from .external_services.pinecone import get_pinecone_index +from .external_services.postgre_vector import get_postgre_cursor from .external_services.openai import create_embedding, query_chat_gpt_forcing_json +import numpy as np def build_slack_message_link(workspace_name, channel_id, message_timestamp, thread_timestamp=None): @@ -40,29 +42,35 @@ def smart_query(namespace, query, username: str): logging.info(f"Smart Query: embedding created in {round(create_embedding_time, 2)}s, " f"trace_id = {trace_id}") - pinecone_search_start_time = time.perf_counter() - query_results = get_pinecone_index().query( - queries=[query_vector], - top_k=50, - namespace=namespace, - include_values=False, - includeMetadata=True - ) - query_matches = query_results['results'][0]['matches'] - pinecone_search_time = time.perf_counter() - pinecone_search_start_time - logging.info(f"Smart Query: Pinecone search finished in {round(pinecone_search_time, 2)}s, " + db_search_start_time = time.perf_counter() + get_postgre_cursor().execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) + query_matches = get_postgre_cursor().fetchall() + + # query( + # queries=[query_vector], + # top_k=50, + # namespace=namespace, + # include_values=False, + # includeMetadata=True + # ) + # query_matches = query_results['results'][0]['matches'] + db_search_time = time.perf_counter() - db_search_start_time + logging.info(f"Smart Query: Postgre search finished in {round(db_search_time, 2)}s, " f"trace_id = {trace_id}") gpt_request_start_time = time.perf_counter() - messages_for_gpt = [ - { - "id": qm["id"], - "text": qm["metadata"]["text_without_context"] - if "text_without_context" in qm["metadata"] - else qm["metadata"]["text"] - } - for qm in query_matches - ] + messages_for_gpt = [] + for qm in query_matches: + metadata = json.loads(qm[3]) + messages_for_gpt.append( + { + "id": qm[2], + "text": metadata["text_without_context"] + if "text_without_context" in metadata + else metadata["text"] + } + ) + prompt = (f"Act as a Smart Search Engine that can logically infer an answer to the given query. " f"Be aware of today's date: {str(date.today())} and use it in your conclusions.\n\n" "Here is a list of Slack messages in JSON:\n" @@ -104,7 +112,7 @@ def smart_query(namespace, query, username: str): ) used_messages = list( filter( - lambda match: match["id"] + lambda match: match[2] in result["messages"], query_matches ) ) diff --git a/src/services/api_service.py b/src/services/api_service.py index d99e0dc..99d16ae 100644 --- a/src/services/api_service.py +++ b/src/services/api_service.py @@ -91,6 +91,8 @@ def get_team_subscription(team_id): # @todo cache results def is_smart_search_available(team_id): + # if STANDALONE: + # return True subscription = get_team_subscription(team_id) return subscription["semantic_search_enabled"] is True From 848d68cc04a4f9ff41adb5e26d1f062fcaadbd78 Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Thu, 7 Dec 2023 10:00:12 +0800 Subject: [PATCH 02/10] pgvector --- src/services/api_service.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/services/api_service.py b/src/services/api_service.py index 99d16ae..096c3ac 100644 --- a/src/services/api_service.py +++ b/src/services/api_service.py @@ -92,7 +92,7 @@ def get_team_subscription(team_id): # @todo cache results def is_smart_search_available(team_id): # if STANDALONE: - # return True + # return True subscription = get_team_subscription(team_id) return subscription["semantic_search_enabled"] is True From 7d903aaeb53f8c45aa6e5291ef3eac01b5f66ca9 Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 07:54:17 +0800 Subject: [PATCH 03/10] upgrade minor change for the coding style. --- .env.example | 12 ++++++------ README.md | 8 ++++++++ requirements.txt | 2 +- src/semantic_search/semantic_search/config.py | 10 +++++----- src/semantic_search/semantic_search/query.py | 16 +++++----------- src/services/api_service.py | 2 -- 6 files changed, 25 insertions(+), 25 deletions(-) diff --git a/.env.example b/.env.example index 3625b30..672b14f 100644 --- a/.env.example +++ b/.env.example @@ -23,9 +23,9 @@ PINECONE_INDEX=index-name PINECONE_ENVIRONMENT=environment PINECONE_KEY=key -# POSTGRE -POSTGRE_HOST=localhost -POSTGRE_PORT=5432 -POSTGRE_DATABASE=db-name -POSTGRE_USER=postgres -POSTGRE_PASSWORD=super-secret \ No newline at end of file +# POSTGRES +POSTGRES_HOST=localhost +POSTGRES_PORT=5432 +POSTGRES_DATABASE=db-name +POSTGRES_USER=postgres +POSTGRES_PASSWORD=super-secret \ No newline at end of file diff --git a/README.md b/README.md index 76b1da7..ec95772 100644 --- a/README.md +++ b/README.md @@ -132,6 +132,14 @@ API_BASE_URL=not-needed-for-standalone LOG_LEVEL=DEBUG STANDALONE=true SLACK_USER_ID=U01JZQZQZQZ # Put a your workspace admin user ID if you know it + +# POSTGRES +POSTGRES_HOST=localhost +POSTGRES_PORT=5432 +POSTGRES_DATABASE=db-name +POSTGRES_USER=postgres +POSTGRES_PASSWORD=super-secret + ``` - Update SLACK_BOT_TOKEN (OAuth token), SLACK_SIGNING_SECRET, OPENAI_API_KEY ([Click here to learn how to get an API key from OpenAI](https://www.maisieai.com/help/how-to-get-an-openai-api-key-for-chatgpt)), and SLACK_USER_ID ([Click here how to get your Slack user ID](https://www.workast.com/help/article/how-to-find-a-slack-user-id/)) - Have venv installed `python3 -m pip install virtualenv` diff --git a/requirements.txt b/requirements.txt index c3628e7..9b951e1 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,4 +5,4 @@ openai==0.28 gunicorn==20.1.0 replicate==0.18.1 psycopg[binary] -psycopg2-binary \ No newline at end of file +psycopg2-binary diff --git a/src/semantic_search/semantic_search/config.py b/src/semantic_search/semantic_search/config.py index 9c077bf..8c80c27 100644 --- a/src/semantic_search/semantic_search/config.py +++ b/src/semantic_search/semantic_search/config.py @@ -56,16 +56,16 @@ def is_standalone() -> bool: return os.environ.get('STANDALONE') == 'true' def get_postgre_host()-> str : - return os.environ.get('POSTGRE_HOST') + return os.environ.get('POSTGRES_HOST') def get_postgre_port()-> str : - return os.environ.get('POSTGRE_PORT') + return os.environ.get('POSTGRES_PORT') def get_postgre_database()-> str : - return os.environ.get('POSTGRE_DATABASE') + return os.environ.get('POSTGRES_DATABASE') def get_postgre_user()-> str : - return os.environ.get('POSTGRE_USER') + return os.environ.get('POSTGRES_USER') def get_postgre_password()-> str : - return os.environ.get('POSTGRE_PASSWORD') + return os.environ.get('POSTGRES_PASSWORD') diff --git a/src/semantic_search/semantic_search/query.py b/src/semantic_search/semantic_search/query.py index 017e490..71c86aa 100644 --- a/src/semantic_search/semantic_search/query.py +++ b/src/semantic_search/semantic_search/query.py @@ -8,6 +8,8 @@ from .external_services.openai import create_embedding, query_chat_gpt_forcing_json import numpy as np +CHUNK_ID = 2 +METADATA = 3 def build_slack_message_link(workspace_name, channel_id, message_timestamp, thread_timestamp=None): base_url = f"https://{workspace_name}.slack.com/archives/{channel_id}/" @@ -46,14 +48,6 @@ def smart_query(namespace, query, username: str): get_postgre_cursor().execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) query_matches = get_postgre_cursor().fetchall() - # query( - # queries=[query_vector], - # top_k=50, - # namespace=namespace, - # include_values=False, - # includeMetadata=True - # ) - # query_matches = query_results['results'][0]['matches'] db_search_time = time.perf_counter() - db_search_start_time logging.info(f"Smart Query: Postgre search finished in {round(db_search_time, 2)}s, " f"trace_id = {trace_id}") @@ -61,10 +55,10 @@ def smart_query(namespace, query, username: str): gpt_request_start_time = time.perf_counter() messages_for_gpt = [] for qm in query_matches: - metadata = json.loads(qm[3]) + metadata = json.loads(qm[METADATA]) messages_for_gpt.append( { - "id": qm[2], + "id": qm[CHUNK_ID], "text": metadata["text_without_context"] if "text_without_context" in metadata else metadata["text"] @@ -112,7 +106,7 @@ def smart_query(namespace, query, username: str): ) used_messages = list( filter( - lambda match: match[2] + lambda match: match[CHUNK_ID] in result["messages"], query_matches ) ) diff --git a/src/services/api_service.py b/src/services/api_service.py index 096c3ac..d99e0dc 100644 --- a/src/services/api_service.py +++ b/src/services/api_service.py @@ -91,8 +91,6 @@ def get_team_subscription(team_id): # @todo cache results def is_smart_search_available(team_id): - # if STANDALONE: - # return True subscription = get_team_subscription(team_id) return subscription["semantic_search_enabled"] is True From 9ff1d54ae4449a161f7350e70eebc276742cd0ba Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 08:29:01 +0800 Subject: [PATCH 04/10] fix spelling issue and python library version --- requirements.txt | 4 +-- src/semantic_search/semantic_search/config.py | 10 +++--- .../{postgre_vector.py => postgres_vector.py} | 20 ++++++------ .../handle_indexation_tasks.py | 4 +-- .../semantic_search/load_messages.py | 32 +++++++++---------- src/semantic_search/semantic_search/query.py | 6 ++-- 6 files changed, 38 insertions(+), 38 deletions(-) rename src/semantic_search/semantic_search/external_services/{postgre_vector.py => postgres_vector.py} (57%) diff --git a/requirements.txt b/requirements.txt index 9b951e1..8079dc4 100644 --- a/requirements.txt +++ b/requirements.txt @@ -4,5 +4,5 @@ slack-bolt==1.18.0 openai==0.28 gunicorn==20.1.0 replicate==0.18.1 -psycopg[binary] -psycopg2-binary +psycopg[binary]==3.1.14 +psycopg2-binary==2.9.9 diff --git a/src/semantic_search/semantic_search/config.py b/src/semantic_search/semantic_search/config.py index 8c80c27..cb4f473 100644 --- a/src/semantic_search/semantic_search/config.py +++ b/src/semantic_search/semantic_search/config.py @@ -55,17 +55,17 @@ def get_api_shared_secret() -> str: def is_standalone() -> bool: return os.environ.get('STANDALONE') == 'true' -def get_postgre_host()-> str : +def get_postgres_host()-> str : return os.environ.get('POSTGRES_HOST') -def get_postgre_port()-> str : +def get_postgres_port()-> str : return os.environ.get('POSTGRES_PORT') -def get_postgre_database()-> str : +def get_postgres_database()-> str : return os.environ.get('POSTGRES_DATABASE') -def get_postgre_user()-> str : +def get_postgres_user()-> str : return os.environ.get('POSTGRES_USER') -def get_postgre_password()-> str : +def get_postgres_password()-> str : return os.environ.get('POSTGRES_PASSWORD') diff --git a/src/semantic_search/semantic_search/external_services/postgre_vector.py b/src/semantic_search/semantic_search/external_services/postgres_vector.py similarity index 57% rename from src/semantic_search/semantic_search/external_services/postgre_vector.py rename to src/semantic_search/semantic_search/external_services/postgres_vector.py index 1280cdb..ca363bd 100644 --- a/src/semantic_search/semantic_search/external_services/postgre_vector.py +++ b/src/semantic_search/semantic_search/external_services/postgres_vector.py @@ -1,16 +1,16 @@ from pgvector.psycopg2 import register_vector import psycopg2 -from ..config import get_postgre_host, get_postgre_port, get_postgre_database, get_postgre_user, get_postgre_password +from ..config import get_postgres_host, get_postgres_port, get_postgres_database, get_postgres_user, get_postgres_password conn = None try: conn = psycopg2.connect( - host=get_postgre_host(), - database=get_postgre_database(), - user=get_postgre_user(), - password=get_postgre_password(), - port=get_postgre_port()) + host=get_postgres_host(), + database=get_postgres_database(), + user=get_postgres_user(), + password=get_postgres_password(), + port=get_postgres_port()) cur = conn.cursor() @@ -25,14 +25,14 @@ print(error) -def get_postgre_cursor(): +def get_postgres_cursor(): return cur -def postgre_commit(): +def postgres_commit(): if conn is not None: conn.commit() -def postgre_excute(postgre_cursor): +def postgres_excute(postgres_cursor): if conn is not None: - postgre_cursor.excute(postgre_cursor) \ No newline at end of file + postgres_cursor.excute(postgres_cursor) \ No newline at end of file diff --git a/src/semantic_search/semantic_search/handle_indexation_tasks.py b/src/semantic_search/semantic_search/handle_indexation_tasks.py index ff0bfb1..b05d272 100644 --- a/src/semantic_search/semantic_search/handle_indexation_tasks.py +++ b/src/semantic_search/semantic_search/handle_indexation_tasks.py @@ -4,7 +4,7 @@ from flask import Flask, request, jsonify from .config import get_google_tasks_service_account from .external_services.pinecone import get_pinecone_index -from .external_services.postgre_vector import get_postgre_cursor +from .external_services.postgres_vector import get_postgres_cursor from .google_tasks import queue_task from .load_messages import index_messages from .external_services.slack_api import load_previous_messages_with_pointer @@ -44,7 +44,7 @@ def handle_task(): [messages, next_last_message, start_from] = load_previous_messages_with_pointer(namespace, channel_id, last_message_id, BULK_SIZE) logging.info(f"Task: {task_id}, Iteration Number: {iteration_number}") logging.info(f"Task: {task_id}, Number of Actual Messages: {len(messages)}") - index_messages(channel_id, messages, start_from, get_postgre_cursor(), namespace) + index_messages(channel_id, messages, start_from, get_postgres_cursor(), namespace) if next_last_message is not None: queue_task({ 'task_id': task_id, diff --git a/src/semantic_search/semantic_search/load_messages.py b/src/semantic_search/semantic_search/load_messages.py index 465e96d..0271e52 100644 --- a/src/semantic_search/semantic_search/load_messages.py +++ b/src/semantic_search/semantic_search/load_messages.py @@ -4,7 +4,7 @@ from .config import CONTEXT_LENGTH from .external_services.pinecone import get_pinecone_index -from .external_services.postgre_vector import get_postgre_cursor, postgre_commit +from .external_services.postgres_vector import get_postgres_cursor, postgres_commit from .external_services.openai import create_embeddings, summarize_thread_with_chat_gpt_3_5 import datetime from .external_services.slack_api import fetch_thread_messages, fetch_channel_messages, is_thread, \ @@ -121,7 +121,7 @@ def attach_header(embeddings: List[Embedding], header: Embedding) -> List[Embedd return part_with_header + [embedding.add_header(header) for embedding in part_without_header] -def index_messages(channel_id, messages, start_from, postgre_cursor, namespace): +def index_messages(channel_id, messages, start_from, postgres_cursor, namespace): total_messages = len(messages) logging.info("Replacing User IDs with User Names in the messages") @@ -167,7 +167,7 @@ def index_messages(channel_id, messages, start_from, postgre_cursor, namespace): messages_for_embedding = list(filter(lambda emb_t: len(emb_t.text) != 0, messages_for_embedding)) logging.info(f"Removed empty messages, {str(len(messages_for_embedding))} messages left") - insert_db_embeddings(messages_for_embedding, postgre_cursor, namespace) + insert_db_embeddings(messages_for_embedding, postgres_cursor, namespace) def index_whole_channel(namespace, channel_id): @@ -179,10 +179,10 @@ def index_whole_channel(namespace, channel_id): total_messages = len(messages) logging.info(f"Filtering out service messages, left {str(total_messages)} messages") - index_messages(channel_id, messages, 0, get_postgre_cursor(), namespace) + index_messages(channel_id, messages, 0, get_postgres_cursor(), namespace) -def insert_db_embeddings(messages_for_embedding: List[Embedding], postgre_cursor, namespace): +def insert_db_embeddings(messages_for_embedding: List[Embedding], postgres_cursor, namespace): logging.info("Starting embeddings creation for the generated messages") chunk_size = 30 # for OpenAI embedding_chunks = [messages_for_embedding[i:i + chunk_size] for i in @@ -196,19 +196,19 @@ def insert_db_embeddings(messages_for_embedding: List[Embedding], postgre_cursor for i in range(len(chunk)): metadata = json.dumps(chunk[i].to_metadata()) - postgre_cursor.execute('INSERT INTO embedding (namespace, chunk_id, metadata, values) VALUES (%s, %s, %s, %s)', (namespace, chunk[i].id, metadata, embeddings[i])) + postgres_cursor.execute('INSERT INTO embedding (namespace, chunk_id, metadata, values) VALUES (%s, %s, %s, %s)', (namespace, chunk[i].id, metadata, embeddings[i])) - postgre_commit() + postgres_commit() except: logging.exception("Couldn't insert embeddings") -def delete_db_embedding(embeddings: List[Embedding], postgre_cursor, namespace): +def delete_db_embedding(embeddings: List[Embedding], postgres_cursor, namespace): ids = list(map(lambda emb: emb.id, embeddings)) logging.info(f"Deleting embeddings for {str(ids)}") ids_to_delete_str = ", ".join(map(str, ids)) - postgre_cursor.execute('DELETE FROM embedding WHERE chunk_id IN (%s) AND namespace=%s',(ids_to_delete_str, namespace)) - postgre_commit() + postgres_cursor.execute('DELETE FROM embedding WHERE chunk_id IN (%s) AND namespace=%s',(ids_to_delete_str, namespace)) + postgres_commit() def handle_message_update_and_reindex(body): event = body['event'] @@ -221,15 +221,15 @@ def handle_message_update_and_reindex(body): if not is_actual_message(message): return embedding = slack_message_to_embedding(channel_id, message) - delete_db_embedding([embedding], get_postgre_cursor(), team_id) + delete_db_embedding([embedding], get_postgres_cursor(), team_id) if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgre_cursor(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgres_cursor(), team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgre_cursor(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgres_cursor(), team_id) return if 'subtype' in event and event['subtype'] == 'message_changed': # processing a message update @@ -239,12 +239,12 @@ def handle_message_update_and_reindex(body): return if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgre_cursor(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgres_cursor(), team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH)[1:] # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgre_cursor(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgres_cursor(), team_id) return if 'subtype' not in event: message = event @@ -259,7 +259,7 @@ def handle_message_update_and_reindex(body): embeddings = generate_embedding_for_message(team_id, channel_id, message_id, thread_ts) insert_db_embeddings( embeddings, - get_postgre_cursor(), + get_postgres_cursor(), team_id ) diff --git a/src/semantic_search/semantic_search/query.py b/src/semantic_search/semantic_search/query.py index 71c86aa..ef8ff61 100644 --- a/src/semantic_search/semantic_search/query.py +++ b/src/semantic_search/semantic_search/query.py @@ -4,7 +4,7 @@ import uuid from datetime import date from .external_services.pinecone import get_pinecone_index -from .external_services.postgre_vector import get_postgre_cursor +from .external_services.postgres_vector import get_postgres_cursor from .external_services.openai import create_embedding, query_chat_gpt_forcing_json import numpy as np @@ -45,8 +45,8 @@ def smart_query(namespace, query, username: str): f"trace_id = {trace_id}") db_search_start_time = time.perf_counter() - get_postgre_cursor().execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) - query_matches = get_postgre_cursor().fetchall() + get_postgres_cursor().execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) + query_matches = get_postgres_cursor().fetchall() db_search_time = time.perf_counter() - db_search_start_time logging.info(f"Smart Query: Postgre search finished in {round(db_search_time, 2)}s, " From 0dde0050165521969fc1b678996e07b972d2f8b1 Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 10:06:59 +0800 Subject: [PATCH 05/10] update readme about installing pgvector extension --- .env.example | 5 ----- README.md | 3 +++ 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/.env.example b/.env.example index 672b14f..824ca0c 100644 --- a/.env.example +++ b/.env.example @@ -18,11 +18,6 @@ STANDALONE=true SLACK_USER_ID=U01JZQZQZQZ USE_FALLBACK=false -# PINECONE -PINECONE_INDEX=index-name -PINECONE_ENVIRONMENT=environment -PINECONE_KEY=key - # POSTGRES POSTGRES_HOST=localhost POSTGRES_PORT=5432 diff --git a/README.md b/README.md index ec95772..65d1c58 100644 --- a/README.md +++ b/README.md @@ -110,6 +110,9 @@ Scroll down further and use [Haly Profile Image](https://github.com/UpMortem/sla 5. After installing, you will find a Bot user OAuth token. Save this for later use. ![image](https://github.com/UpMortem/slack-bot/assets/5354324/7d866eee-a7e6-4059-b422-bae8ac9016a3) +*** +6. Finally, Install PostgreSQL and Postgres vector extention using pgvector. Install[PostgreSQL download](https://www.postgresql.org/download/) Install [vector extension](https://github.com/pgvector/pgvector#installation) or (https://github.com/pgvector/pgvector#installation) + ### Configure your project - In a terminal git clone the project. [See the offical documentation if you do not know how to do this](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository). - cd in the slack-bot directory From 88f10e11a4b4342d97f77ffb13dca2fd933c56be Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 10:09:06 +0800 Subject: [PATCH 06/10] minor change for the readme --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 65d1c58..1666573 100644 --- a/README.md +++ b/README.md @@ -111,7 +111,7 @@ Scroll down further and use [Haly Profile Image](https://github.com/UpMortem/sla ![image](https://github.com/UpMortem/slack-bot/assets/5354324/7d866eee-a7e6-4059-b422-bae8ac9016a3) *** -6. Finally, Install PostgreSQL and Postgres vector extention using pgvector. Install[PostgreSQL download](https://www.postgresql.org/download/) Install [vector extension](https://github.com/pgvector/pgvector#installation) or (https://github.com/pgvector/pgvector#installation) +6. Finally, Install PostgreSQL and Postgres vector extention using pgvector. Install [PostgreSQL download](https://www.postgresql.org/download/) Install [vector extension](https://github.com/pgvector/pgvector#installation) or (https://github.com/pgvector/pgvector#installation) ### Configure your project - In a terminal git clone the project. [See the offical documentation if you do not know how to do this](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository). From 34d380b96de85a9372307becbd87d4b94354e873 Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 10:15:53 +0800 Subject: [PATCH 07/10] minor change again --- README.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 1666573..8bbbf78 100644 --- a/README.md +++ b/README.md @@ -111,7 +111,9 @@ Scroll down further and use [Haly Profile Image](https://github.com/UpMortem/sla ![image](https://github.com/UpMortem/slack-bot/assets/5354324/7d866eee-a7e6-4059-b422-bae8ac9016a3) *** -6. Finally, Install PostgreSQL and Postgres vector extention using pgvector. Install [PostgreSQL download](https://www.postgresql.org/download/) Install [vector extension](https://github.com/pgvector/pgvector#installation) or (https://github.com/pgvector/pgvector#installation) +6. Finally, Install [PostgreSQL](https://www.postgresql.org/download/) and Postgres vector extention using [pgvector](https://github.com/pgvector/pgvector#installation). Additional [Install Methods](https://github.com/pgvector/pgvector#installation) for pgvector. + +*** ### Configure your project - In a terminal git clone the project. [See the offical documentation if you do not know how to do this](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository). From 850ace758d67319bc8abc9618ee52ad4ae3fbf3b Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Fri, 8 Dec 2023 10:18:26 +0800 Subject: [PATCH 08/10] changed readme finally --- README.md | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/README.md b/README.md index 8bbbf78..18265ff 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,7 @@ https://haly.ai 1. You must have a Slack organization where either you or an administrator can approve a new application. 2. The ability to git clone a repo and run commands in either a Windows, Mac, or Linux terminal. 3. Install [python](https://www.python.org/downloads/) and [pip](https://pip.pypa.io/en/stable/installation/). +4. Finally, Install [PostgreSQL](https://www.postgresql.org/download/) and Postgres vector extention using [pgvector](https://github.com/pgvector/pgvector#installation). Additional [Install Methods](https://github.com/pgvector/pgvector#installation) for pgvector. ### Create your Slack bot: 1. Go to https://api.slack.com/apps and hit the "Create New App" green button. Select "From an app manifest" option. @@ -111,9 +112,7 @@ Scroll down further and use [Haly Profile Image](https://github.com/UpMortem/sla ![image](https://github.com/UpMortem/slack-bot/assets/5354324/7d866eee-a7e6-4059-b422-bae8ac9016a3) *** -6. Finally, Install [PostgreSQL](https://www.postgresql.org/download/) and Postgres vector extention using [pgvector](https://github.com/pgvector/pgvector#installation). Additional [Install Methods](https://github.com/pgvector/pgvector#installation) for pgvector. -*** ### Configure your project - In a terminal git clone the project. [See the offical documentation if you do not know how to do this](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository). From 799e160c5d145bcd22e84b96cb4318d28a1db0ae Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Mon, 11 Dec 2023 17:09:17 +0800 Subject: [PATCH 09/10] abstract layer for the vector database --- README.md | 3 -- .../external_services/pinecone.py | 9 ---- .../external_services/postgres_vector.py | 38 -------------- .../external_services/slack_api.py | 5 +- .../vector_databases/__init__.py | 0 .../vector_databases/pinecone.py | 42 ++++++++++++++++ .../vector_databases/postgres.py | 50 +++++++++++++++++++ .../vector_databases/vector_database.py | 14 ++++++ .../vector_databases/vector_instance.py | 17 +++++++ .../handle_indexation_tasks.py | 4 +- .../semantic_search/load_messages.py | 38 +++++--------- src/semantic_search/semantic_search/query.py | 42 ++++++++-------- 12 files changed, 161 insertions(+), 101 deletions(-) delete mode 100644 src/semantic_search/semantic_search/external_services/pinecone.py delete mode 100644 src/semantic_search/semantic_search/external_services/postgres_vector.py create mode 100644 src/semantic_search/semantic_search/external_services/vector_databases/__init__.py create mode 100644 src/semantic_search/semantic_search/external_services/vector_databases/pinecone.py create mode 100644 src/semantic_search/semantic_search/external_services/vector_databases/postgres.py create mode 100644 src/semantic_search/semantic_search/external_services/vector_databases/vector_database.py create mode 100644 src/semantic_search/semantic_search/external_services/vector_databases/vector_instance.py diff --git a/README.md b/README.md index 18265ff..d549d69 100644 --- a/README.md +++ b/README.md @@ -111,9 +111,6 @@ Scroll down further and use [Haly Profile Image](https://github.com/UpMortem/sla 5. After installing, you will find a Bot user OAuth token. Save this for later use. ![image](https://github.com/UpMortem/slack-bot/assets/5354324/7d866eee-a7e6-4059-b422-bae8ac9016a3) -*** - - ### Configure your project - In a terminal git clone the project. [See the offical documentation if you do not know how to do this](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository). - cd in the slack-bot directory diff --git a/src/semantic_search/semantic_search/external_services/pinecone.py b/src/semantic_search/semantic_search/external_services/pinecone.py deleted file mode 100644 index 6b930a2..0000000 --- a/src/semantic_search/semantic_search/external_services/pinecone.py +++ /dev/null @@ -1,9 +0,0 @@ -import pinecone - -from ..config import get_pinecone_key, get_pinecone_environment, get_pinecone_index_name - -pinecone.init(api_key=get_pinecone_key(), environment=get_pinecone_environment()) - - -def get_pinecone_index() -> 'pinecone.Index': - return pinecone.Index(get_pinecone_index_name()) diff --git a/src/semantic_search/semantic_search/external_services/postgres_vector.py b/src/semantic_search/semantic_search/external_services/postgres_vector.py deleted file mode 100644 index ca363bd..0000000 --- a/src/semantic_search/semantic_search/external_services/postgres_vector.py +++ /dev/null @@ -1,38 +0,0 @@ -from pgvector.psycopg2 import register_vector -import psycopg2 - -from ..config import get_postgres_host, get_postgres_port, get_postgres_database, get_postgres_user, get_postgres_password - -conn = None -try: - conn = psycopg2.connect( - host=get_postgres_host(), - database=get_postgres_database(), - user=get_postgres_user(), - password=get_postgres_password(), - port=get_postgres_port()) - - cur = conn.cursor() - - cur.execute('CREATE EXTENSION IF NOT EXISTS vector') - register_vector(cur) - - # cur.execute('DROP TABLE IF EXISTS embedding') - cur.execute('CREATE TABLE IF NOT EXISTS embedding (id bigserial PRIMARY KEY, namespace text, chunk_id text, metadata text, values vector)') - - conn.commit() -except (Exception, psycopg2.DatabaseError) as error: - print(error) - - -def get_postgres_cursor(): - return cur - - -def postgres_commit(): - if conn is not None: - conn.commit() - -def postgres_excute(postgres_cursor): - if conn is not None: - postgres_cursor.excute(postgres_cursor) \ No newline at end of file diff --git a/src/semantic_search/semantic_search/external_services/slack_api.py b/src/semantic_search/semantic_search/external_services/slack_api.py index 294cefe..a0b7f73 100644 --- a/src/semantic_search/semantic_search/external_services/slack_api.py +++ b/src/semantic_search/semantic_search/external_services/slack_api.py @@ -114,9 +114,8 @@ def slack_names_map(team_id): def load_previous_messages(team_id: str, channel_id: str, last_message_id: str, number: int): - result = load_previous_messages_with_pointer(team_id, channel_id, last_message_id, number) - messages = result[0][-number:] - return messages + (messages, _, _) = load_previous_messages_with_pointer(team_id, channel_id, last_message_id, number) + return messages[-number:] def load_subsequent_messages(team_id: str, channel_id: str, first_message_id: str, number: int): diff --git a/src/semantic_search/semantic_search/external_services/vector_databases/__init__.py b/src/semantic_search/semantic_search/external_services/vector_databases/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/semantic_search/semantic_search/external_services/vector_databases/pinecone.py b/src/semantic_search/semantic_search/external_services/vector_databases/pinecone.py new file mode 100644 index 0000000..81ee04d --- /dev/null +++ b/src/semantic_search/semantic_search/external_services/vector_databases/pinecone.py @@ -0,0 +1,42 @@ +import pinecone +from .vector_database import VectorDatabase +from ...config import get_pinecone_key, get_pinecone_environment, get_pinecone_index_name + +class Pinecone(VectorDatabase): + def __init__(self): + pinecone.init(api_key=get_pinecone_key(), environment=get_pinecone_environment()) + + # overriding abstract method + def insert(self, embeddings, chunk, namespace): + items = [] + + for i in range(len(chunk)): + items.append({ + 'id': chunk[i].id, + 'values': embeddings[i], + 'metadata': chunk[i].to_metadata() + }) + + self.get_pinecone_index().upsert( + vectors=items, + namespace=namespace + ) + + # overriding abstract method + def delete(self, ids, namespace): + self.get_pinecone_index().delete(ids=ids, namespace=namespace) + + # overriding abstract method + def select(self, query_vector, namespace): + query_results = self.get_pinecone_index().query( + queries=[query_vector], + top_k=50, + namespace=namespace, + include_values=False, + includeMetadata=True + ) + return query_results['results'][0]['matches'] + + def get_pinecone_index() -> 'pinecone.Index': + return pinecone.Index(get_pinecone_index_name()) + diff --git a/src/semantic_search/semantic_search/external_services/vector_databases/postgres.py b/src/semantic_search/semantic_search/external_services/vector_databases/postgres.py new file mode 100644 index 0000000..12b2242 --- /dev/null +++ b/src/semantic_search/semantic_search/external_services/vector_databases/postgres.py @@ -0,0 +1,50 @@ +from pgvector.psycopg2 import register_vector +import psycopg2 +from .vector_database import VectorDatabase +import json +import numpy as np + +from ...config import get_postgres_host, get_postgres_port, get_postgres_database, get_postgres_user, get_postgres_password + +class Postgres(VectorDatabase): + def __init__(self): + try: + self.conn = psycopg2.connect( + host=get_postgres_host(), + database=get_postgres_database(), + user=get_postgres_user(), + password=get_postgres_password(), + port=get_postgres_port()) + + self.cur = self.conn.cursor() + + self.cur.execute('CREATE EXTENSION IF NOT EXISTS vector') + register_vector(self.cur) + + self.cur.execute('CREATE TABLE IF NOT EXISTS embedding (id bigserial PRIMARY KEY, namespace text, chunk_id text, metadata text, values vector)') + self.conn.commit() + + except (Exception, psycopg2.DatabaseError) as error: + print(error) + + # overriding abstract method + def insert(self, embeddings, chunk, namespace): + for i in range(len(chunk)): + metadata = json.dumps(chunk[i].to_metadata()) + self.cur.execute('INSERT INTO embedding (namespace, chunk_id, metadata, values) VALUES (%s, %s, %s, %s)', (namespace, chunk[i].id, metadata, embeddings[i])) + + self.conn.commit() + + # overriding abstract method + def delete(self, ids, namespace): + self.cur.execute('DELETE FROM embedding WHERE chunk_id IN (%s) AND namespace=%s',(ids, namespace)) + self.conn.commit() + + # overriding abstract method + def select(self, query_vector): + self.cur.execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) + embeddings = self.cur.fetchall() + output = [dict(id=chunk_id, metadata=json.loads(metadata)) for id, namespace, chunk_id, metadata, values in embeddings] + return output + + diff --git a/src/semantic_search/semantic_search/external_services/vector_databases/vector_database.py b/src/semantic_search/semantic_search/external_services/vector_databases/vector_database.py new file mode 100644 index 0000000..7b25884 --- /dev/null +++ b/src/semantic_search/semantic_search/external_services/vector_databases/vector_database.py @@ -0,0 +1,14 @@ +from abc import ABC, abstractmethod + +class VectorDatabase(ABC): + @abstractmethod + def insert(self, embeddings, chunk, namespace): + pass + + @abstractmethod + def delete(self, ids, namespace): + pass + + @abstractmethod + def select(self, query_vector): + pass \ No newline at end of file diff --git a/src/semantic_search/semantic_search/external_services/vector_databases/vector_instance.py b/src/semantic_search/semantic_search/external_services/vector_databases/vector_instance.py new file mode 100644 index 0000000..59af463 --- /dev/null +++ b/src/semantic_search/semantic_search/external_services/vector_databases/vector_instance.py @@ -0,0 +1,17 @@ +from ...config import get_postgres_host, get_pinecone_environment +from .postgres import Postgres +from .pinecone import Pinecone + +postgres_instance = None +pinecone_instance = None + +def get_db_instance(): + global postgres_instance, pinecone_instance + if get_postgres_host(): + if postgres_instance is None: + postgres_instance = Postgres() + return postgres_instance + else: + if pinecone_instance is None: + pinecone_instance = Pinecone() + return pinecone_instance \ No newline at end of file diff --git a/src/semantic_search/semantic_search/handle_indexation_tasks.py b/src/semantic_search/semantic_search/handle_indexation_tasks.py index b05d272..eb38fd3 100644 --- a/src/semantic_search/semantic_search/handle_indexation_tasks.py +++ b/src/semantic_search/semantic_search/handle_indexation_tasks.py @@ -3,8 +3,6 @@ import os from flask import Flask, request, jsonify from .config import get_google_tasks_service_account -from .external_services.pinecone import get_pinecone_index -from .external_services.postgres_vector import get_postgres_cursor from .google_tasks import queue_task from .load_messages import index_messages from .external_services.slack_api import load_previous_messages_with_pointer @@ -44,7 +42,7 @@ def handle_task(): [messages, next_last_message, start_from] = load_previous_messages_with_pointer(namespace, channel_id, last_message_id, BULK_SIZE) logging.info(f"Task: {task_id}, Iteration Number: {iteration_number}") logging.info(f"Task: {task_id}, Number of Actual Messages: {len(messages)}") - index_messages(channel_id, messages, start_from, get_postgres_cursor(), namespace) + index_messages(channel_id, messages, start_from, namespace) if next_last_message is not None: queue_task({ 'task_id': task_id, diff --git a/src/semantic_search/semantic_search/load_messages.py b/src/semantic_search/semantic_search/load_messages.py index 0271e52..67ed4a1 100644 --- a/src/semantic_search/semantic_search/load_messages.py +++ b/src/semantic_search/semantic_search/load_messages.py @@ -1,18 +1,13 @@ import logging import copy from typing import List, Dict - from .config import CONTEXT_LENGTH -from .external_services.pinecone import get_pinecone_index -from .external_services.postgres_vector import get_postgres_cursor, postgres_commit from .external_services.openai import create_embeddings, summarize_thread_with_chat_gpt_3_5 import datetime from .external_services.slack_api import fetch_thread_messages, fetch_channel_messages, is_thread, \ is_actual_message, \ slack_names_map, filter_messages, load_previous_messages, load_subsequent_messages - -import numpy as np -import json +from .external_services.vector_databases.vector_instance import get_db_instance class Embedding: def __init__(self, channel_id, id, text, ts, thread_ts=None, author_id=None): @@ -121,7 +116,7 @@ def attach_header(embeddings: List[Embedding], header: Embedding) -> List[Embedd return part_with_header + [embedding.add_header(header) for embedding in part_without_header] -def index_messages(channel_id, messages, start_from, postgres_cursor, namespace): +def index_messages(channel_id, messages, start_from, namespace): total_messages = len(messages) logging.info("Replacing User IDs with User Names in the messages") @@ -167,7 +162,7 @@ def index_messages(channel_id, messages, start_from, postgres_cursor, namespace) messages_for_embedding = list(filter(lambda emb_t: len(emb_t.text) != 0, messages_for_embedding)) logging.info(f"Removed empty messages, {str(len(messages_for_embedding))} messages left") - insert_db_embeddings(messages_for_embedding, postgres_cursor, namespace) + insert_db_embeddings(messages_for_embedding, namespace) def index_whole_channel(namespace, channel_id): @@ -179,10 +174,10 @@ def index_whole_channel(namespace, channel_id): total_messages = len(messages) logging.info(f"Filtering out service messages, left {str(total_messages)} messages") - index_messages(channel_id, messages, 0, get_postgres_cursor(), namespace) + index_messages(channel_id, messages, 0, namespace) -def insert_db_embeddings(messages_for_embedding: List[Embedding], postgres_cursor, namespace): +def insert_db_embeddings(messages_for_embedding: List[Embedding], namespace): logging.info("Starting embeddings creation for the generated messages") chunk_size = 30 # for OpenAI embedding_chunks = [messages_for_embedding[i:i + chunk_size] for i in @@ -193,22 +188,16 @@ def insert_db_embeddings(messages_for_embedding: List[Embedding], postgres_curso counter += len(chunk) try: embeddings = create_embeddings([embedding_message.text for embedding_message in chunk]) - - for i in range(len(chunk)): - metadata = json.dumps(chunk[i].to_metadata()) - postgres_cursor.execute('INSERT INTO embedding (namespace, chunk_id, metadata, values) VALUES (%s, %s, %s, %s)', (namespace, chunk[i].id, metadata, embeddings[i])) - - postgres_commit() + get_db_instance().insert(embeddings, chunk, namespace) except: logging.exception("Couldn't insert embeddings") -def delete_db_embedding(embeddings: List[Embedding], postgres_cursor, namespace): +def delete_db_embedding(embeddings: List[Embedding], namespace): ids = list(map(lambda emb: emb.id, embeddings)) logging.info(f"Deleting embeddings for {str(ids)}") ids_to_delete_str = ", ".join(map(str, ids)) - postgres_cursor.execute('DELETE FROM embedding WHERE chunk_id IN (%s) AND namespace=%s',(ids_to_delete_str, namespace)) - postgres_commit() + get_db_instance().delete(ids_to_delete_str, namespace) def handle_message_update_and_reindex(body): event = body['event'] @@ -221,15 +210,15 @@ def handle_message_update_and_reindex(body): if not is_actual_message(message): return embedding = slack_message_to_embedding(channel_id, message) - delete_db_embedding([embedding], get_postgres_cursor(), team_id) + delete_db_embedding([embedding], team_id) if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgres_cursor(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH - 1) # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgres_cursor(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, team_id) return if 'subtype' in event and event['subtype'] == 'message_changed': # processing a message update @@ -239,12 +228,12 @@ def handle_message_update_and_reindex(body): return if message.get('thread_ts') is not None: # just reindex the whole thread - index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, get_postgres_cursor(), team_id) + index_messages(channel_id, load_previous_messages(team_id, channel_id, message.get('thread_ts'), 1), 0, team_id) return message_ts = message['ts'] messages_for_reindex = load_previous_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH) + load_subsequent_messages(team_id, channel_id, message_ts, CONTEXT_LENGTH)[1:] # reindex surrounding messages - index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, get_postgres_cursor(), team_id) + index_messages(channel_id, messages_for_reindex, CONTEXT_LENGTH - 1, team_id) return if 'subtype' not in event: message = event @@ -259,7 +248,6 @@ def handle_message_update_and_reindex(body): embeddings = generate_embedding_for_message(team_id, channel_id, message_id, thread_ts) insert_db_embeddings( embeddings, - get_postgres_cursor(), team_id ) diff --git a/src/semantic_search/semantic_search/query.py b/src/semantic_search/semantic_search/query.py index ef8ff61..52931ef 100644 --- a/src/semantic_search/semantic_search/query.py +++ b/src/semantic_search/semantic_search/query.py @@ -3,13 +3,8 @@ import time import uuid from datetime import date -from .external_services.pinecone import get_pinecone_index -from .external_services.postgres_vector import get_postgres_cursor from .external_services.openai import create_embedding, query_chat_gpt_forcing_json -import numpy as np - -CHUNK_ID = 2 -METADATA = 3 +from .external_services.vector_databases.vector_instance import get_db_instance def build_slack_message_link(workspace_name, channel_id, message_timestamp, thread_timestamp=None): base_url = f"https://{workspace_name}.slack.com/archives/{channel_id}/" @@ -45,8 +40,7 @@ def smart_query(namespace, query, username: str): f"trace_id = {trace_id}") db_search_start_time = time.perf_counter() - get_postgres_cursor().execute('SELECT * FROM embedding ORDER BY values <-> %s LIMIT 50', (np.array(query_vector),)) - query_matches = get_postgres_cursor().fetchall() + query_matches = get_db_instance().select(query_vector) db_search_time = time.perf_counter() - db_search_start_time logging.info(f"Smart Query: Postgre search finished in {round(db_search_time, 2)}s, " @@ -54,17 +48,25 @@ def smart_query(namespace, query, username: str): gpt_request_start_time = time.perf_counter() messages_for_gpt = [] - for qm in query_matches: - metadata = json.loads(qm[METADATA]) - messages_for_gpt.append( - { - "id": qm[CHUNK_ID], - "text": metadata["text_without_context"] - if "text_without_context" in metadata - else metadata["text"] - } - ) - + # for qm in query_matches: + # metadata = json.loads(qm[METADATA]) + # messages_for_gpt.append( + # { + # "id": qm[CHUNK_ID], + # "text": metadata["text_without_context"] + # if "text_without_context" in metadata + # else metadata["text"] + # } + # ) + messages_for_gpt = [ + { + "id": qm["id"], + "text": qm["metadata"]["text_without_context"] + if "text_without_context" in qm["metadata"] + else qm["metadata"]["text"] + } + for qm in query_matches + ] prompt = (f"Act as a Smart Search Engine that can logically infer an answer to the given query. " f"Be aware of today's date: {str(date.today())} and use it in your conclusions.\n\n" "Here is a list of Slack messages in JSON:\n" @@ -106,7 +108,7 @@ def smart_query(namespace, query, username: str): ) used_messages = list( filter( - lambda match: match[CHUNK_ID] + lambda match: match["id"] in result["messages"], query_matches ) ) From 7f2e58fb221e6fd2bf828cf5bf1cab490a32daae Mon Sep 17 00:00:00 2001 From: tozzen0121 Date: Tue, 12 Dec 2023 08:09:29 +0800 Subject: [PATCH 10/10] remove dead code --- src/semantic_search/semantic_search/query.py | 11 ----------- 1 file changed, 11 deletions(-) diff --git a/src/semantic_search/semantic_search/query.py b/src/semantic_search/semantic_search/query.py index 52931ef..620bb7b 100644 --- a/src/semantic_search/semantic_search/query.py +++ b/src/semantic_search/semantic_search/query.py @@ -47,17 +47,6 @@ def smart_query(namespace, query, username: str): f"trace_id = {trace_id}") gpt_request_start_time = time.perf_counter() - messages_for_gpt = [] - # for qm in query_matches: - # metadata = json.loads(qm[METADATA]) - # messages_for_gpt.append( - # { - # "id": qm[CHUNK_ID], - # "text": metadata["text_without_context"] - # if "text_without_context" in metadata - # else metadata["text"] - # } - # ) messages_for_gpt = [ { "id": qm["id"],