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'

Streaming Projections Report - {report_date.strftime("%B %d, %Y")}

') html_parts.append(f'

PM: {pm_name}

') html_parts.append(f'

{len(today_df)} track(s) overperforming today.

') table_html = '' table_html += '' for col in ['ISRC', 'Track Name', 'Artist', 'Expected', 'Actual', '% Diff', 'Spikes in Last 7D']: table_html += f'' table_html += '' for _, row in today_df.iterrows(): pct_diff = (row['ERROR'] / row['EXPECTED_STREAMS'] * 100) if row['EXPECTED_STREAMS'] > 0 else 0 isrc_link = f'https://insights.awal.com/song/{row["ISRC"]}/?dimension=all&potGraphType=LINE' table_html += '' table_html += f'' table_html += f'' table_html += f'' table_html += f'' table_html += f'' table_html += f'' table_html += f'' table_html += '' table_html += '
{col}
{row["ISRC"]}{row["TRACK_NAME"]}{row["ARTIST_NAME"]}{int(row["EXPECTED_STREAMS"]):,}{int(row["ACTUAL_STREAMS"]):,}+{pct_diff:.0f}%{int(row["SPIKE_DAYS"])}
' html_parts.append(table_html) html_parts.append('
') inline_images = [] for chart_idx, (_, row) in enumerate(today_df.iterrows()): isrc = row['ISRC'] track_name = row['TRACK_NAME'] artist_name = row['ARTIST_NAME'] actual = int(row['ACTUAL_STREAMS']) expected = int(row['EXPECTED_STREAMS']) pct_diff = (row['ERROR'] / row['EXPECTED_STREAMS'] * 100) if row['EXPECTED_STREAMS'] > 0 else 0 spike_days = int(row['SPIKE_DAYS']) isrc_link = f'https://insights.awal.com/song/{isrc}/?dimension=all&potGraphType=LINE' track_df = pm_df[pm_df['ISRC'] == isrc].sort_values('DATE').reset_index(drop=True) fig, ax = plt.subplots(figsize=(8, 3.5)) # Expected range band upper = track_df['EXPECTED_STREAMS'] + track_df['NORMAL_ERROR_STD'] * 2 lower = (track_df['EXPECTED_STREAMS'] - track_df['NORMAL_ERROR_STD'] * 2).clip(lower=0) ax.fill_between(track_df['DATE'], lower, upper, alpha=0.15, color='gray', label='Expected Range') # Expected streams line ax.plot(track_df['DATE'], track_df['EXPECTED_STREAMS'], color='gray', linestyle='--', linewidth=1.5, label='Expected Streams') # Actual streams line ax.plot(track_df['DATE'], track_df['ACTUAL_STREAMS'], color='#1f77b4', linewidth=2, marker='o', markersize=5, label='Actual Streams') # Red dots for overperforming days overperf = track_df[track_df['IS_OVERPERFORMING'] == 1] if not overperf.empty: ax.scatter(overperf['DATE'], overperf['ACTUAL_STREAMS'], color='red', s=30, zorder=5) # Formatting ax.set_title(f"{track_name} - {artist_name}", fontsize=11, fontweight='bold', loc='left') ax.set_xlabel('Date', fontsize=9) ax.set_ylabel('Streams', fontsize=9) ax.yaxis.set_major_formatter(mticker.FuncFormatter(lambda x, p: f'{int(x):,}')) ax.xaxis.set_major_formatter(mdates.DateFormatter('%b %d')) ax.xaxis.set_major_locator(mdates.DayLocator(interval=2)) plt.xticks(rotation=45, fontsize=8) plt.yticks(fontsize=8) ax.spines['top'].set_visible(False) ax.spines['right'].set_visible(False) ax.grid(axis='y', alpha=0.3) if chart_idx == 0: ax.legend(fontsize=8, loc='upper left', framealpha=0.9) plt.tight_layout() # Convert to base64 PNG buf = BytesIO() fig.savefig(buf, format='png', dpi=150, bbox_inches='tight') plt.close(fig) buf.seek(0) # img_b64 = base64.b64encode(buf.read()).decode('utf-8') summary = ( f'

{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'') html_parts.append('
') full_html = f''' {chr(10).join(html_parts)} ''' # output_path = os.path.join(os.getcwd(), 'report_preview.html') # with open(output_path, 'w') as f: # f.write(full_html) print(f'Report generated: {len(today_df)} tracks for {pm_name}') # print(f'Saved to {output_path}') print(f'Preparing email') email_subject = f'AWAL Insights: Roster Alert {today}' full_email = full_html destination_emails = destination_email_list['EMAILS'].tolist() email_outcome = prep_and_send_email(email_params,destination_emails,email_subject,full_email,inline_images) if email_outcome: logger.info('Email sent.') else: logger.info('Error: Email did not send')