import pandas as pd import numpy as np import streamlit as st import pytz import datetime as dt from dateutil.relativedelta import relativedelta, FR import time import snowflake.connector import os import re import boto3 import matplotlib.pyplot as plt import matplotlib.dates as mdates import matplotlib.ticker as mticker # import base64 from io import BytesIO def init_secrets(): session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name='us-east-1' ) # sf_user = st.secrets['SNOWFLAKE_USER'] # sf_account = st.secrets['SNOWFLAKE_ACCOUNT'] # sf_warehouse = st.secrets['SNOWFLAKE_WAREHOUSE'] # sf_password = st.secrets['PEM_KEY_PASSWORD'] # sf_pem_key = st.secrets['PEM_KEY'] sf_user = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_USER")["SecretString"] sf_account = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_ACCOUNT")["SecretString"] sf_warehouse = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_WAREHOUSE")["SecretString"] sf_password = client.get_secret_value(SecretId="dev/awal-ar/PEM_KEY_PASSWORD")["SecretString"] sf_pem_key = client.get_secret_value(SecretId="dev/awal-ar/PEM_KEY")["SecretString"] from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives import serialization p_key = serialization.load_pem_private_key( sf_pem_key.encode('utf-8').decode('unicode_escape').encode("utf-8"), password=sf_password.encode('utf-8'), backend=default_backend() ) pkb = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption()) sf_secrets = { 'sf_user':sf_user, 'sf_account':sf_account, 'sf_warehouse':sf_warehouse, 'sf_password':sf_password, 'pkb':pkb } return sf_secrets def orcd_query(sql_query,sf_params): ctx = sf_params['snowflake'].connector.connect( user=sf_params['sf_user'], private_key=sf_params['pkb'], account=sf_params['sf_account'], warehouse=sf_params['sf_warehouse'] ) cs = ctx.cursor(sf_params['snowflake'].connector.DictCursor) try: cs.execute(sql_query) result = cs.fetchall() finally: cs.close() ctx.close() result = pd.DataFrame(result) return result def setup_email_params(): send_the_email = True send_only_to_myself = False client = boto3.client('ses',region_name='us-east-1') email_params = { 'client':client, 'send_the_email':send_the_email, 'send_only_to_myself':send_only_to_myself } return email_params def prep_and_send_email(email_params,destination_emails,email_subject,full_email,inline_images): if email_params['send_the_email']: from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from email.mime.image import MIMEImage if email_params['send_only_to_myself']: destination_emails = ['joselyn.ho@awal.com']#,'erik.gunnarsson@awal.com'] msg = MIMEMultipart("related") msg["Subject"] = email_subject msg["From"] = 'awalinsights@dev.theorchard.io' msg["To"] = ", ".join(destination_emails) reply_to_emails = ['joselyn.ho@awal.com','erik.gunnarsson@awal.com','insights@awal.com'] msg["Reply-To"] = ", ".join(reply_to_emails) # Alternative part (HTML body) alt = MIMEMultipart("alternative") msg.attach(alt) alt.attach(MIMEText("Please view this email in an HTML-capable client.", "plain")) alt.attach( MIMEText(full_email, "html", "utf-8") ) # Attach inline images for img in inline_images: cid = img["cid"] img_buf = img["img_buf"] mime_img = MIMEImage(img_buf.getvalue(),_subtype="png") # CRITICAL: angle brackets are required # mime_img.add_header("Content-ID", f"<{cid}>") # mime_img.add_header("Content-Disposition", "inline", filename=path) mime_img.add_header("Content-ID", f"<{cid}>") mime_img.add_header( "Content-Disposition", "inline", filename=f"{cid}.png" ) # mime_img["Content-ID"] = f"<{cid}>" # mime_img["Content-Disposition"] = f'inline; filename="{cid}.png"' msg.attach(mime_img) raw_email = msg.as_bytes() print(f"Email size: {len(raw_email)/1024:.1f} KB") # Build SES params params = { "Source": 'awalinsights@dev.theorchard.io', "Destinations": destination_emails, "RawMessage": { "Data": msg.as_bytes() }, "SourceArn": 'arn:aws:ses:us-east-1:103233932089:identity/dev.theorchard.io', } return email_params['client'].send_raw_email(**params) # return response['ResponseMetadata']['HTTPStatusCode'] == 200 else: # this is for local testing where no email is sent. Just opens a browser with the contents import webbrowser import os # Save the file file_path = os.path.abspath("preview_email.html") with open(file_path, "w", encoding="utf-8") as f: f.write(full_email) # Open it in the default web browser webbrowser.open(f"file://{file_path}") def setup_logger(): import logging logger = logging.getLogger() logger.setLevel(logging.INFO) # Add a logger handler if none exists if not logger.hasHandlers(): handler = logging.StreamHandler() formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) logger.addHandler(handler) return logger ################################################### ################################################### def handler(event,context): logger = setup_logger() sf_params = init_secrets() sf_params['snowflake'] = snowflake email_params = setup_email_params() destination_email_list = orcd_query(""" select * from awal.awal_ar.email_alerts where CATEGORY='ROSTER_ALERTS' """,sf_params) projections_df = orcd_query(""" select * from dev_engineering.egunnarsson.daily_streaming_alerts """,sf_params) et_timezone = pytz.timezone('US/Eastern') today = dt.datetime.now(et_timezone).date() df = projections_df.copy() df['DATE'] = pd.to_datetime(df['DATE']) pm_name = 'Phil McDonald' pm_df = df[df['PRODUCT_MANAGER'] == pm_name].copy() if pm_df.empty: print(f'No flagged tracks for {pm_name} today.') print(f'Available PMs with flagged tracks: {df["PRODUCT_MANAGER"].unique().tolist()}') else: report_date = pm_df['DATE'].max() today_df = pm_df[pm_df['DATE'] == report_date].copy() today_df = today_df.sort_values(['ARTIST_NAME', 'ACTUAL_STREAMS'], ascending=[True, False]) html_parts = [] html_parts.append(f'
{len(today_df)} track(s) overperforming today.
') table_html = '| {col} | ' table_html += '||||||
|---|---|---|---|---|---|---|
| {row["ISRC"]} | ' table_html += f'{row["TRACK_NAME"]} | ' table_html += f'{row["ARTIST_NAME"]} | ' table_html += f'{int(row["EXPECTED_STREAMS"]):,} | ' table_html += f'{int(row["ACTUAL_STREAMS"]):,} | ' table_html += f'+{pct_diff:.0f}% | ' table_html += f'{int(row["SPIKE_DAYS"])} | ' table_html += '
{track_name} overperformed expectations on ' f'{report_date.strftime("%B %d, %Y")}. ' f'Actual streams were {pct_diff:.0f}% more than expected. ' f'It has spiked {spike_days} day(s) in the last 7 days.
' f'Visit Insights for more.
' ) safe_isrc = re.sub(r'[^A-Za-z0-9_-]', '_', str(isrc)) cid = f'chart_{safe_isrc}' inline_images.append({ 'cid': cid, 'img_buf': buf }) html_parts.append(summary) html_parts.append(f'