| import gradio as gr |
| import os |
| import psycopg2 |
| import json |
| import shopify |
| import threading |
| import select |
| import psycopg2.extensions |
| from datetime import datetime |
| from openai import OpenAI |
| from dotenv import load_dotenv |
| from hubspot import HubSpot |
| from orchestrator import Orchestrator |
| from ucp_wrapper import ucp_get_order_details |
| from shopify_tools import get_product_details |
| from research_tools import web_search, summarize_content, analyze_product_trends |
| from fulfillment_tools import fulfill_shopify_order |
| from product_builder_tools import create_shopify_product |
| from email_tools import list_emails, send_email |
| from auditor_tools import check_margin, log_audit_event |
| from supplier_tools import evaluate_supplier, log_evaluation |
|
|
| |
| load_dotenv() |
|
|
| |
| openai_client = OpenAI(api_key=os.getenv("OPENAI_API_KEY")) |
| hubspot = HubSpot(access_token=os.getenv("HUBSPOT_ACCESS_TOKEN", "")) |
| shopify_session = shopify.Session(os.getenv("SHOPIFY_SHOP_URL", ""), "2024-04", os.getenv("SHOPIFY_ACCESS_TOKEN", "")) |
| shopify.ShopifyResource.activate_session(shopify_session) |
| orchestrator = Orchestrator() |
|
|
| |
| tool_map = { |
| "ucp_get_order_details": ucp_get_order_details, |
| "get_product_details": get_product_details, |
| "web_search": web_search, |
| "summarize_content": summarize_content, |
| "analyze_product_trends": analyze_product_trends, |
| "fulfill_shopify_order": fulfill_shopify_order, |
| "create_shopify_product": create_shopify_product, |
| "list_emails": list_emails, |
| "send_email": send_email, |
| "check_margin": check_margin, |
| "log_audit_event": log_audit_event, |
| "evaluate_supplier": evaluate_supplier, |
| "log_evaluation": log_evaluation |
| } |
|
|
| |
| class DatabaseManager: |
| def __init__(self): |
| self.db_url = os.getenv("DATABASE_URL") |
| self._init_db() |
|
|
| def _get_conn(self): |
| return psycopg2.connect(self.db_url) |
|
|
| def _init_db(self): |
| conn = self._get_conn() |
| cursor = conn.cursor() |
| cursor.execute("CREATE TABLE IF NOT EXISTS sessions (id SERIAL PRIMARY KEY, name TEXT, created_at TIMESTAMP)") |
| cursor.execute("""CREATE TABLE IF NOT EXISTS messages |
| (id SERIAL PRIMARY KEY, session_id INTEGER, role TEXT, content TEXT, timestamp TIMESTAMP)""") |
| cursor.execute("""CREATE TABLE IF NOT EXISTS research_reports |
| (id SERIAL PRIMARY KEY, session_id INTEGER, topic TEXT, content TEXT, timestamp TIMESTAMP)""") |
| |
| cursor.execute("""CREATE TABLE IF NOT EXISTS transaction_logs ( |
| id SERIAL PRIMARY KEY, session_id INTEGER, agent_name TEXT, tool_name TEXT, |
| input_data JSONB, output_data JSONB, is_compliant BOOLEAN DEFAULT TRUE, timestamp TIMESTAMP DEFAULT NOW())""") |
| conn.commit() |
| conn.close() |
|
|
| def create_session(self, name): |
| conn = self._get_conn() |
| cursor = conn.cursor() |
| cursor.execute("INSERT INTO sessions (name, created_at) VALUES (%s, %s) RETURNING id", (name, datetime.now())) |
| session_id = cursor.fetchone()[0] |
| conn.commit() |
| conn.close() |
| return session_id |
|
|
| def get_sessions(self): |
| conn = self._get_conn() |
| cursor = conn.cursor() |
| cursor.execute("SELECT id, name FROM sessions ORDER BY created_at DESC") |
| sessions = cursor.fetchall() |
| conn.close() |
| return sessions |
|
|
| def save_message(self, session_id, role, content): |
| conn = self._get_conn() |
| cursor = conn.cursor() |
| cursor.execute("INSERT INTO messages (session_id, role, content, timestamp) VALUES (%s, %s, %s, %s)", |
| (session_id, role, content, datetime.now())) |
| conn.commit() |
| conn.close() |
|
|
| def get_history(self, session_id): |
| if not session_id: return [] |
| conn = self._get_conn() |
| cursor = conn.cursor() |
| cursor.execute("SELECT role, content FROM messages WHERE session_id = %s ORDER BY timestamp ASC", (session_id,)) |
| rows = cursor.fetchall() |
| conn.close() |
| return [(r[1] if r[0] == 'user' else None, r[1] if r[0] == 'assistant' else None) for r in rows] |
|
|
| db = DatabaseManager() |
|
|
| |
| def listen_for_webhooks(): |
| conn = psycopg2.connect(os.getenv("DATABASE_URL")) |
| conn.set_isolation_level(psycopg2.extensions.ISOLATION_LEVEL_AUTOCOMMIT) |
| cursor = conn.cursor() |
| cursor.execute("LISTEN new_webhook;") |
| |
| print("Agent listening for webhooks...") |
| while True: |
| if select.select([conn], [], [], 5) == ([], [], []): |
| continue |
| conn.poll() |
| while conn.notifies: |
| notify = conn.notifies.pop(0) |
| event_data = json.loads(notify.payload) |
| print(f"New Webhook Received: {event_data['event_type']}") |
| |
| |
|
|
| |
| threading.Thread(target=listen_for_webhooks, daemon=True).start() |
|
|
| |
| def render_storefront_component(component_type, product_handle): |
| store_domain = os.getenv("SHOPIFY_SHOP_URL", "") |
| return f""" |
| <script src="https://cdn.shopify.com/s/assets/storefront/load-components.js"></script> |
| <shopify-context shop-url="{store_domain}"></shopify-context> |
| <{component_type} handle="{product_handle}"></{component_type}> |
| """ |
|
|
| |
| def chat_wrapper(message, history, session_id, model): |
| if not message or not session_id: return "", history |
| db.save_message(session_id, "user", message) |
| |
| |
| agent_key = orchestrator.get_agent_for_intent(message) |
| sys_prompt, tools = orchestrator.get_agent_config(agent_key) |
| |
| messages = [{"role": "system", "content": sys_prompt}] + [{"role": h[0] and "user" or "assistant", "content": h[0] or h[1]} for h in history if h[0] or h[1]] |
| messages.append({"role": "user", "content": message}) |
| |
| |
| response = openai_client.chat.completions.create(model=model, messages=messages, tools=tools) |
| msg_obj = response.choices[0].message |
| |
| |
| if msg_obj.tool_calls: |
| messages.append(msg_obj) |
| for tool_call in msg_obj.tool_calls: |
| func_name = tool_call.function.name |
| func_args = json.loads(tool_call.function.arguments) |
| |
| func_res = tool_map[func_name](**func_args) |
| |
| |
| |
| log_audit_event(session_id, agent_key, func_name, True) |
| |
| messages.append({ |
| "role": "tool", |
| "tool_call_id": tool_call.id, |
| "content": str(func_res) |
| }) |
| |
| response = openai_client.chat.completions.create(model=model, messages=messages) |
| assistant_reply = response.choices[0].message.content |
| else: |
| assistant_reply = msg_obj.content |
|
|
| db.save_message(session_id, "assistant", assistant_reply) |
| history.append((message, assistant_reply)) |
| return "", history |
|
|
| |
| with gr.Blocks(theme=gr.themes.Soft()) as demo: |
| gr.Markdown("# 🚀 Optimized AI Agent Dashboard (Auditor Edition)") |
| chatbot = gr.Chatbot(height=400) |
| |
| with gr.Tabs(): |
| with gr.TabItem("Session Manager"): |
| session_name = gr.Textbox(label="New Session Name") |
| create_btn = gr.Button("Create Session") |
| session_dropdown = gr.Dropdown(label="Select Session", choices=db.get_sessions()) |
| create_btn.click(lambda name: db.create_session(name), [session_name], [session_dropdown]) |
| |
| with gr.TabItem("Chat"): |
| msg = gr.Textbox(label="Input") |
| model = gr.Dropdown(choices=["gpt-3.5-turbo", "gpt-4o", "gpt-4-turbo"], value="gpt-3.5-turbo") |
| msg.submit(chat_wrapper, [msg, chatbot, session_dropdown, model], [msg, chatbot]) |
| session_dropdown.change(lambda sid: db.get_history(sid), [session_dropdown], [chatbot]) |
|
|
| with gr.TabItem("Storefront"): |
| comp_type = gr.Dropdown(choices=["shopify-product-card", "shopify-buy-button"], label="Component Type") |
| handle = gr.Textbox(label="Product Handle") |
| render_btn = gr.Button("Render Component") |
| html_out = gr.HTML() |
| render_btn.click(render_storefront_component, [comp_type, handle], [html_out]) |
|
|
| if __name__ == "__main__": |
| demo.launch() |
|
|