Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions backend/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,15 @@ class Settings(BaseSettings):
sentry_dsn: str = ""
project_id: str = "localmate"
environment: str = "prod"
# Square — global token deprecated for production (per-client OAuth is
# authoritative via square_oauth.get_valid_token). Kept for sandbox/dev fallback.
square_access_token: str = ""
square_environment: str = "sandbox"
square_app_id: str = ""
square_app_secret: str = ""
square_oauth_redirect_path: str = "/auth/square-callback"
square_webhook_signature_key: str = ""
menu_images_bucket: str = "menu-images"
supabase_jwt_secret: str = ""

# --- Phase 0: queue / worker / billing-portal infra ---
Expand All @@ -35,6 +42,11 @@ class Settings(BaseSettings):
dashboard_url: str = "" # Stripe portal return_url base
stripe_portal_config_id: str = "" # Stripe portal configuration id (bpc_...)

# --- Phase 4: GBP Pub/Sub provisioning (D15-B full automation) ---
gcp_project_id: str = "" # GCP project hosting the Pub/Sub topic
gcp_sa_json: str = "" # service-account JSON key (topic-admin) for Pub/Sub REST
gbp_pubsub_topic_name: str = "gbp-reviews" # shared topic short name

class Config:
env_file = ".env.local"
env_file_encoding = "utf-8"
Expand Down
98 changes: 79 additions & 19 deletions backend/jobs/competitor_watch.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,24 @@

from db import get_db
from services.claude import generate_competitor_brief
from services.structured_extract import (
extract_structured,
detect_prices_from_text,
diff_structured,
)
from utils.retry import retry_on_failure

logger = logging.getLogger(__name__)


@retry_on_failure()
async def snapshot_website(url: str) -> tuple[str, str]:
"""Fetch a competitor URL, strip non-content elements, and return (md5_hash, clean_text).
async def snapshot_website(url: str) -> tuple[str, str, str]:
"""Fetch a competitor URL, strip non-content elements, and return (md5_hash, clean_text, raw_html).

Returns ("", "") when the key is missing or the response is empty.
The raw HTML is returned so the caller can run ``extract_structured`` on
the original (with JSON-LD ``<script>`` blocks intact, before stripping).

Returns ("", "", "") when the key is missing or the response is empty.
"""
async with httpx.AsyncClient() as client:
resp = await client.get(
Expand All @@ -26,45 +34,82 @@ async def snapshot_website(url: str) -> tuple[str, str]:
timeout=15,
)
resp.raise_for_status()
soup = BeautifulSoup(resp.text, "html.parser")
raw_html = resp.text
soup = BeautifulSoup(raw_html, "html.parser")
for tag in soup(["script", "style", "nav", "footer"]):
tag.decompose()
clean_text = " ".join(soup.get_text().split())
md5_hex = hashlib.md5(clean_text.encode()).hexdigest()
return md5_hex, clean_text
return md5_hex, clean_text, raw_html


def _build_structured(raw_html: str, clean_text: str) -> dict:
"""Extract structured data from HTML, with regex text fallback for prices."""
structured = extract_structured(raw_html)
# Regex fallback: supplement prices when JSON-LD found none.
if not structured["prices"]:
text_prices = detect_prices_from_text(clean_text)
if text_prices:
structured["prices"] = text_prices
return structured


def _format_structured_diff(diff: list[dict]) -> list[str]:
"""Format structured diffs into human-readable lines for the brief."""
lines = []
for d in diff:
kind = d["kind"]
name = d["name"]
if kind == "changed":
lines.append(f"{name} price ${d['old']} → ${d['new']}")
elif kind == "added":
lines.append(f"{name} added at ${d['new']}")
elif kind == "removed":
lines.append(f"{name} removed (was ${d['old']})")
return lines


async def detect_changes(client_id: str, competitor_url: str) -> dict | None:
"""Compare the latest snapshot hash against a fresh fetch.
"""Compare the latest snapshot against a fresh fetch.

Returns None when the content is unchanged or the fetch failed.
Returns a dict with change details on difference, including the new snapshot row id.

The caller is responsible for calling ``run_competitor_snapshots_all_clients``
which invokes this per client/URL pair.
Returns a dict with change details on difference, including:
- ``structured_diff``: field-level diffs from JSON-LD / price extraction
- ``snapshot_id``: the new snapshot row id
- ``prev_text`` / ``curr_text``: text snippets (fallback for the brief)
"""
db = get_db()

last_resp = (
db.table("competitor_snapshots")
.select("content_hash, content_text")
.select("content_hash, content_text, structured_data")
.eq("client_id", client_id)
.eq("competitor_url", competitor_url)
.order("captured_at", desc=True)
.limit(1)
.execute()
)

new_hash, new_text = await snapshot_website(competitor_url)
new_hash, new_text, raw_html = await snapshot_website(competitor_url)
if not new_hash:
return None

new_structured = _build_structured(raw_html, new_text)

last_text = ""
last_structured: dict = {}
structured_diff: list[dict] = []
if last_resp.data:
last_row = last_resp.data[0]
if last_row["content_hash"] == new_hash:
return None
last_text = last_row.get("content_text") or ""
last_structured = last_row.get("structured_data") or {}
try:
structured_diff = diff_structured(last_structured, new_structured)
except Exception as e:
logger.warning("diff_structured failed for %s: %s", competitor_url, e)
structured_diff = []

insert_resp = (
db.table("competitor_snapshots")
Expand All @@ -73,6 +118,8 @@ async def detect_changes(client_id: str, competitor_url: str) -> dict | None:
"competitor_url": competitor_url,
"content_hash": new_hash,
"content_text": new_text,
"structured_data": new_structured,
"structured_diff": structured_diff,
})
.execute()
)
Expand All @@ -84,6 +131,7 @@ async def detect_changes(client_id: str, competitor_url: str) -> dict | None:
"prev_text": last_text,
"curr_text": new_text,
"snapshot_id": new_id,
"structured_diff": structured_diff,
}


Expand All @@ -107,8 +155,11 @@ async def run_competitor_snapshots_all_clients() -> None:

1. Snapshot each URL listed in ``competitor_urls``.
2. If changes are detected, collect them into a list.
3. Call ``generate_competitor_brief`` once per changed client.
4. Persist the brief text and extracted threat level on each new snapshot row.
3. Build ``changes_summary`` from structured diffs first (concrete
field/value changes), falling back to text snippets when no structured
signal exists.
4. Call ``generate_competitor_brief`` once per changed client.
5. Persist the brief text and extracted threat level on each new snapshot row.

Each client is wrapped in its own try/except so one failure never crashes
the entire job run.
Expand Down Expand Up @@ -153,11 +204,20 @@ async def run_competitor_snapshots_all_clients() -> None:

parts = []
for c in changes:
parts.append(
f"Competitor: {c['url']}\n"
f"Previous snippet: {c['prev_text'][:300]}\n"
f"Current snippet: {c['curr_text'][:300]}"
)
# Build summary from structured diffs first (concrete, specific).
structured_lines = _format_structured_diff(c.get("structured_diff", []))
if structured_lines:
parts.append(
f"Competitor: {c['url']}\nStructured changes:\n"
+ "\n".join(structured_lines)
)
else:
# Fall back to text snippets when no structured signal exists.
parts.append(
f"Competitor: {c['url']}\n"
f"Previous snippet: {c['prev_text'][:300]}\n"
f"Current snippet: {c['curr_text'][:300]}"
)
changes_summary = "\n\n".join(parts)

brief = await _generate_brief_safe(business_name, changes_summary)
Expand Down
Loading
Loading