# Multimodal RAG

***This notebook works well with the `Data Science 3.0 Python 3` kernel and `ml.t3.medium` instance type.***

***Run the [`0_data_prep.ipynb`](./0_data_prep.ipynb) notebook prior to running this notebook.***

1. We download a subset of data from the [Amazon Berkley Objects](https://amazon-berkeley-objects.s3.amazonaws.com/index.html) dataset. The data includes Amazon products with metadata and catalog images. The metadata includes multiple tags that provide short text description of the product in the image. The data is filtered to only keep images that are associated for tags with description in a given language (`enUS` in our example), to limit the size of the data.

1. We convert all downloaded images into Base64 encoding.

1. An image and its associated text are converted into embeddings in a single `invoke_model` call to the `amazon.titan-embed-image-v1` model. We embed all images in our dataset in this way.

1. These embeddings are then ingested into in-memory [FAISS](https://github.com/facebookresearch/faiss) database to store and search for embeddings vectors. In a real-world scenario, you will likely want to use a persistent data store such as the [vector engine for Amazon OpenSearch Service Serverless](https://aws.amazon.com/opensearch-service/serverless-vector-engine/) or the pgvector extension for PostgreSQL.

1. Now for retrieval, we consider the following scenario: a customer is looking for a product and has a text description of the product and, optionally, an image of the product, we do an embeddings based similarity search by converting the text description and the image into embeddings using the `Amazon Titan Embeddings G1 - Image` model and retrieve the most relevant results from the vector database. 

1. We then further refine these search results by creating a text prompt using the description for the retrieved objects and asking the LLM (`Anthropic Claude v2`) to do the following:
    * Reason through the responses based on the customer's description of what they were looking for and then either accept or reject each of the search results
    * Explain the reasoning why each result was accepted or rejected. 
    * The text response generated by the model along with the accepted set of results (images and text) are returned to the customer.

In [3]:
import sys
#!{sys.executable} -m pip install -r requirements.txt


In [4]:
import os
import io
import sys
import json
import glob
import faiss
import boto3
import base64
import logging
import requests
import numpy as np
import pandas as pd
from PIL import Image
from typing import List
import logging
import boto3
# from botocore.auth import SigV4Auth
from faiss import write_index, read_index
from langchain.llms.bedrock import Bedrock
# from botocore.awsrequest import AWSRequest
#from faiss.swigfaiss_avx2 import IndexFlatIP

# logging.basicConfig() #format='[%(asctime)s] p%(process)s {%(filename)s:%(lineno)d} %(levelname)s - %(message)s', level=logging.INFO)
# logger = logging.getLogger(__name__)


In [5]:
# global constants
import os
LISTINGS_FILE: str = os.path.join("listings", "metadata", "listings_0.json")
LANGUAGE_TO_FILTER: str = "en_US"
IMAGE_ID_TO_FNAME_MAPPING_FILE: str = "images.csv"
ABO_S3_BUCKET: str = "amazon-berkeley-objects"
ABO_S3_PREFIX:str = "images/original"
ABO_S3_BUCKET_PREFIX: str = f"s3://{ABO_S3_BUCKET}/{ABO_S3_PREFIX}"
IMAGE_DATASET_FNAME: str = f"aob_{LANGUAGE_TO_FILTER}.csv"
DATA_DIR: str = "data"
IMAGES_DIR: str = os.path.join(DATA_DIR, "images", LANGUAGE_TO_FILTER)
B64_ENCODED_IMAGES_DIR: str = os.path.join(DATA_DIR, "b64_images", LANGUAGE_TO_FILTER)
VECTOR_DB_DIR: str = os.path.join(DATA_DIR, "vectordb", LANGUAGE_TO_FILTER)
SUCCESSFULLY_EMBEDDED_DIR: str = os.path.join(DATA_DIR, "successfully_embedded", LANGUAGE_TO_FILTER)
IMAGE_DATA_W_SUCCESSFUL_EMBEDDINGS_FPATH: str = os.path.join(SUCCESSFULLY_EMBEDDED_DIR, "data.csv")
os.makedirs(DATA_DIR, exist_ok=True)
os.makedirs(IMAGES_DIR, exist_ok=True)
os.makedirs(B64_ENCODED_IMAGES_DIR, exist_ok=True)
os.makedirs(VECTOR_DB_DIR, exist_ok=True)
os.makedirs(SUCCESSFULLY_EMBEDDED_DIR, exist_ok=True)
FMC_URL: str = "https://bedrock-runtime.us-east-1.amazonaws.com"
FMC_MODEL_ID: str = "amazon.titan-embed-image-v1"
CLAUDE_V2_MODEL_ID: str  = "anthropic.claude-v2"
ACCEPT_ENCODING: str = "application/json"
CONTENT_ENCODING: str = "application/json"
VECTORDB_INDEX_FILE: str = f"aob_{LANGUAGE_TO_FILTER}_index"
VECTOR_DB_INDEX_FPATH: str = os.path.join(VECTOR_DB_DIR, VECTORDB_INDEX_FILE)
K: int = 4
N: int = 10000
MAX_IMAGE_HEIGHT: int = 2048
MAX_IMAGE_WIDTH: int = 2048


In [7]:
def display_image(fpath: str) -> None:
    image = Image.open(fpath)
    return image

# convert all the downloaded files into base64 encoding
def encode_image_to_base64(image_file_path: str):
    with open(image_file_path, "rb") as image_file:
        b64_image = base64.b64encode(image_file.read()).decode('utf8')
    return b64_image


In [8]:
import boto3
bedrock = boto3.client(
    service_name="bedrock-runtime", region_name="us-east-1", endpoint_url=FMC_URL
)


In [9]:
import glob
image_b64_file_list = glob.glob(os.path.join(B64_ENCODED_IMAGES_DIR, "*.b64"))
print(f"there are {len(image_b64_file_list)} base64 encoded images in {B64_ENCODED_IMAGES_DIR}")
image_b64_file_list[:10]


there are 0 base64 encoded images in data/b64_images/en_US


[]

In [16]:
import os
import boto3
import asyncio
import logging
import pandas as pd
from typing import Dict, Generator



# download the files from s3 and convert them into base64 encoded images
def download_image_file(row_tuple: Dict, s3_bucket: str, s3_prefix: str, local_images_dir: str):
    s3 = boto3.client('s3')
    _, row = row_tuple
    # print(f"row type={type(row)}")
    path = row.get('path')
    if path is None:
        return
    local_path = os.path.join(local_images_dir, os.path.basename(path))
    key = f"{s3_prefix}/{path}"
    print(f"going to download {s3_bucket}/{key} to {local_path}")
    with open(local_path, 'wb') as f:
        s3.download_fileobj(s3_bucket, key, f)


def download_images(image_count: int, image_data_fname: str, s3_bucket: str, s3_prefix: str, local_images_dir: str):
    image_data = pd.read_csv(image_data_fname)
    image_count = len(image_data) if image_count > len(image_data) else image_count
    _ = asyncio.run(adownload_all_image_files(image_data.sample(n=image_count).iterrows(), s3_bucket, s3_prefix, local_images_dir))

In [17]:
image_dataset = pd.read_csv(IMAGE_DATASET_FNAME)
print(f"there are {len(image_dataset)} base64 encoded images in {IMAGE_DATASET_FNAME} dataset")

# only keep the rows for which we have image files downloaded already
image_dataset['path_b64'] = image_dataset.path.map(lambda x: os.path.join(B64_ENCODED_IMAGES_DIR, f"{os.path.basename(x)}.b64"))
print(image_dataset.head())

image_b64_list_dataframe = pd.DataFrame(image_b64_file_list)
image_b64_list_dataframe.columns = ["path_b64"]
image_dataset = pd.merge(left=image_dataset, right=image_b64_list_dataframe, how="inner", on="path_b64")
image_dataset


FileNotFoundError: [Errno 2] No such file or directory: 'aob_en_US.csv'

In [None]:
from pandas.core.series import Series
import numpy as np
import numpy
from typing import Dict
def get_embeddings(text: str, image: str) -> numpy.ndarray:
   
    # You can specify either text or image or both
    body = json.dumps(
        {
            "inputText": text,
            "inputImage": image
        }
    )
        
    modelId = FMC_MODEL_ID
    accept = ACCEPT_ENCODING
    contentType = CONTENT_ENCODING

    try:
        response = bedrock.invoke_model(
            body=body, modelId=modelId, accept=accept, contentType=contentType
        )
        response_body = json.loads(response.get("body").read())        
        embeddings = np.array([response_body.get("embedding")]).astype(np.float32)        
    except Exception as e:
        logger.error(f"exception while encoding text={text}, image(truncated)={image[:10]}, exception={e}")
        embeddings = None
    return embeddings

    


In [None]:
%%time

import faiss
import numpy as np

index = None
image_dataset_successful_embeddings_only = []
for _, row in image_dataset.iterrows():
    print(f"encoding image={row['path_b64']}, description={row['description']}")
    path_b64 = row['path_b64']
    # MAX image size supported is 2048 * 2048 pixels
    with open(path_b64, "rb") as image_file:
        input_image_b64 = image_file.read().decode('utf-8')        
    input_text = "No description" if row['description'] is np.nan else row['description']

    embeddings = get_embeddings(input_text, input_image_b64)
    if embeddings is None:
        logger.error(f"error creating embeddings for {row}")
        continue
    image_dataset_successful_embeddings_only.append(row)
    if index is None:
        vector_dimension = embeddings.shape[1]
        index = faiss.IndexFlatIP(vector_dimension)
    
    faiss.normalize_L2(embeddings)
    index.add(embeddings)


In [None]:
print(f"successfully ingested {len(image_dataset_successful_embeddings_only)} images and descriptions into the vector db index")
image_dataset_successful_embeddings_only_df = pd.DataFrame(image_dataset_successful_embeddings_only)
image_dataset_successful_embeddings_only_df.to_csv(IMAGE_DATA_W_SUCCESSFUL_EMBEDDINGS_FPATH, index=False)
image_dataset_successful_embeddings_only_df.head()


In [None]:
print(f"going to save vectordb index with {index.ntotal} to {VECTOR_DB_INDEX_FPATH}")
write_index(index, VECTOR_DB_INDEX_FPATH)


In [None]:
print(f"going to load vectordb index with from {VECTOR_DB_INDEX_FPATH}")
index = read_index(VECTOR_DB_INDEX_FPATH)
print(f"there are {index.ntotal} elements in index loaded from {VECTOR_DB_INDEX_FPATH}, index type={type(index)}")


In [None]:
def find_multimodal_match(search_text: str, search_image: str, index: IndexFlatIP, k: int) -> List:
    print(f"search_text={search_text}, search_image(truncated)={search_image[:100]}, index={index}, k={K}")
    search_vector = get_embeddings(search_text, search_image)
    faiss.normalize_L2(search_vector)
    
    distances, ann = index.search(search_vector, k=k)
    matches = {'distances': distances[0], 'ann': ann[0]}
    results = pd.DataFrame(matches).sort_values(by="distances", ascending=False)
    
    return results


In [None]:
def generate_image(prompt: str, negative_prompts: List):
    body = json.dumps({
        "text_prompts": ([{"text": prompt, "weight": 1.0}] + [{"text": negprompt, "weight": -1.0} for negprompt in negative_prompts]),
        "cfg_scale": 10,
        "seed": 20,
        "steps": 50,
        "style_preset": "photographic",
    })
    modelId = "stability.stable-diffusion-xl"
    accept = "application/json"
    contentType = "application/json"

    try:

        response = bedrock.invoke_model(
            body=body, modelId=modelId, accept=accept, contentType=contentType
        )
        response_body = json.loads(response.get("body").read())

        #print(response_body["result"])
        print(f'{response_body.get("artifacts")[0].get("base64")[0:80]}...')

    except botocore.exceptions.ClientError as error:

        if error.response['Error']['Code'] == 'AccessDeniedException':
            print(f"\x1b[41m{error.response['Error']['Message']}\
                    \nTo troubeshoot this issue please refer to the following resources.\
                    \nhttps://docs.aws.amazon.com/IAM/latest/UserGuide/troubleshoot_access-denied.html\
                    \nhttps://docs.aws.amazon.com/bedrock/latest/userguide/security-iam.html\x1b[0m\n")

        else:
            raise error
    base_64_img_str = response_body.get("artifacts")[0].get("base64")
    #image = Image.open(io.BytesIO(base64.decodebytes(bytes(base_64_img_str, "utf-8"))))
    return base_64_img_str


In [None]:
search_text:str = "black or gray colored running shoe for men, trail running, comfortable, wide toebox. Only show PUMA or Nike."
image_prompt:str = "A zoomed in product image of a black or gray colored trail running shoe for men, comfortable with wide toebox."
negative_prompts: List = ["in focus background", "in focus athlete", "front view"]
search_image = generate_image(image_prompt, negative_prompts)
image = Image.open(io.BytesIO(base64.decodebytes(bytes(search_image, "utf-8"))))
display(image)


In [None]:
image_dataset_successful_embeddings_only_df = pd.read_csv(IMAGE_DATA_W_SUCCESSFUL_EMBEDDINGS_FPATH)
image_dataset_successful_embeddings_only_df.head()


In [None]:
matches = find_multimodal_match(search_text, search_image, index, K)
display(matches)
matches_from_dataset = image_dataset_successful_embeddings_only_df.iloc[list(matches.ann), :]
matches_from_dataset


In [None]:
print(f"search text = \"{search_text}\"")

for i, row in matches_from_dataset.iterrows():
    print(f"--------- Match {i} ----------")
    print(row['description'])
    fpath = os.path.join(IMAGES_DIR, os.path.basename(row['path']))
    print(f"image file={fpath}")
    image = display_image(fpath)
    display(image)


In [None]:
inference_modifier = {
    "max_tokens_to_sample": 4096,
    "temperature": 0,
    "top_k": 250,
    "top_p": 1,
    "stop_sequences": ["\n\nHuman"],
}
textgen_llm = Bedrock(
    model_id=CLAUDE_V2_MODEL_ID,
    client=bedrock,
    model_kwargs=inference_modifier,
)


In [None]:
print(f"search text = \"{search_text}\"")
prompt: str = """Human: You are a shopping assistant bot and helping a customer find what they are looking for from an online product catalog. 
The user query and search results from a product catalog are provided below. Some of the search results may not be relevant based on the user query, 
pick the result which exactly match the user's criteria and return their indexes, if you do not find an exact match then return the next most relevant
results but explicitly say that these are most relevant but not exact matches to the user's criteria.

<user_query>
{}
</user_query>

<search_results>
{}
</search_results>

Now based on the information provided above, provide the index of the most relevant results and also say why you choose those results and for the results that
you did not choose say why not. 
Finally, as the last line of your response, write the most relevant indexes as a comma separated list in a line all by itself as:
<relevant_indices></relevant_indices>

Answer:
"""

search_results: str = ""
for i, row in matches_from_dataset.iterrows():
    search_results += f"Index: {i}, Description: {row['description']}\n\n"
prompt = prompt.format(search_text, search_results)
print(prompt)


In [None]:
response = textgen_llm(prompt)
print(response)


In [None]:
index_list_line = response.split("\n")[-1]
print(f"index_list_line={index_list_line}")


In [None]:
import re
found = re.findall('<relevant_indices>(.*)</relevant_indices>', index_list_line)
print(f"found={found}")
if len(found) > 0:
    indices = [int(i.strip()) for i in found[0].split(",")]
    best_matches = image_dataset_successful_embeddings_only_df.iloc[indices, :]
    for i, row in best_matches.iterrows():
        print(f"--- Best match {i} -----")
        print(f"Description = {row['description']}")
        fpath = os.path.join(IMAGES_DIR, os.path.basename(row['path']))
        print(f"image file={fpath}")
        image = display_image(fpath)
        display(image)
else:
    print(f"the model did not find any valid matches in the multimodal results")
