"""Phase 7: Integrations management — destinations, export, test delivery."""
import csv
import io
import json
from flask import Blueprint, Response, abort, jsonify, redirect, request, session, url_for
from app.helpers import _get_site_owner_password_hash
from app.models import (
create_destination,
delete_destination,
get_destination,
get_site_destinations,
get_submission,
get_user_sites,
list_submissions,
parse_site_fields,
update_destination,
)
from app.routes.auth import login_required
from app.services.integrations import FORMATTERS, get_setup_guide
from app.services.webhook import fire_destination
integrations_bp = Blueprint(
"integrations",
__name__,
url_prefix="/sites/<int:site_id>/integrations",
)
# ─── Integration type metadata ────────────────────────────────────────────────
INTEGRATION_INFO = {
"webhook": {
"name": "Webhook",
"icon": "🔗",
"description": "Send to any HTTP endpoint",
},
"google_sheets": {
"name": "Google Sheets",
"icon": "📊",
"description": "Append to a Google Sheet",
},
"slack": {
"name": "Slack",
"icon": "💬",
"description": "Post to a Slack channel",
},
"discord": {
"name": "Discord",
"icon": "🎮",
"description": "Send to a Discord channel",
},
"telegram": {
"name": "Telegram",
"icon": "✈️",
"description": "Send to a Telegram chat",
},
"airtable": {
"name": "Airtable",
"icon": "📋",
"description": "Add records to Airtable",
},
"notion": {
"name": "Notion",
"icon": "📝",
"description": "Append to a Notion database",
},
}
# ─── Helpers ──────────────────────────────────────────────────────────────────
def _resolve_site(site_id):
"""Return site dict if it belongs to the current user, or abort 404."""
if "user_id" not in session:
abort(404)
user_id = session["user_id"]
sites = get_user_sites(user_id)
site = next((s for s in sites if s["id"] == site_id), None)
if not site:
abort(404)
return site
def _resolve_destination(dest_id, site):
"""Return destination dict if it exists for this site, or abort 404."""
dest = get_destination(dest_id)
if not dest or dest["site_id"] != site["id"]:
abort(404)
return dest
def _mask_url(url, max_visible=12):
"""Mask a URL for display — show prefix + '...' + last few chars."""
if not url:
return ""
visible = url[:max_visible]
return f"{visible}…" if len(url) > max_visible else url
# ─── Destination CRUD ─────────────────────────────────────────────────────────
@integrations_bp.route("/", methods=["GET"])
@login_required
def list_destinations(site_id):
"""List all destinations for a site."""
site = _resolve_site(site_id)
destinations = get_site_destinations(site_id)
result = []
for dest in destinations:
info = INTEGRATION_INFO.get(dest["type"], {})
result.append(
{
"id": dest["id"],
"name": dest["name"],
"type": dest["type"],
"type_name": info.get("name", dest["type"]),
"icon": info.get("icon", "🔌"),
"url": dest.get("url", ""),
"url_display": _mask_url(dest.get("url")),
"enabled": dest.get("enabled", True),
"success_count": dest.get("success_count", 0),
"failure_count": dest.get("failure_count", 0),
"last_status": dest.get("last_status"),
"last_delivered_at": dest.get("last_delivered_at"),
"created_at": dest.get("created_at"),
}
)
return jsonify({"site": {"id": site["id"], "name": site["name"]}, "destinations": result})
@integrations_bp.route("/", methods=["POST"])
@login_required
def create_destination_route(site_id):
"""Create a new destination for a site."""
site = _resolve_site(site_id)
data = request.get_json(silent=True) or {}
name = (data.get("name") or "").strip()
dest_type = (data.get("type") or "").strip()
url = (data.get("url") or "").strip()
config = data.get("config")
errors = []
if not name:
errors.append("Destination name is required")
if not dest_type:
errors.append("Destination type is required")
if dest_type not in FORMATTERS:
errors.append(f"Unknown destination type: {dest_type}")
elif dest_type != "telegram" and not url:
errors.append("URL is required")
if errors:
return jsonify({"errors": errors}), 400
# Validate URL safety if provided
if url:
from app.routes.user_sites import validate_webhook_url
is_safe, error_msg = validate_webhook_url(url)
if not is_safe:
return jsonify({"errors": [f"Invalid URL: {error_msg}"]}), 400
dest = create_destination(site_id, name, dest_type, url, config)
info = INTEGRATION_INFO.get(dest_type, {})
return jsonify(
{
"destination": {
"id": dest["id"],
"name": dest["name"],
"type": dest["type"],
"type_name": info.get("name", dest_type),
"icon": info.get("icon", "🔌"),
"url_display": _mask_url(url),
"enabled": dest.get("enabled", True),
}
}
), 201
@integrations_bp.route("/<int:dest_id>", methods=["PUT"])
@login_required
def update_destination_route(site_id, dest_id):
"""Update an existing destination."""
site = _resolve_site(site_id)
_resolve_destination(dest_id, site)
data = request.get_json(silent=True) or {}
# Validate URL safety if provided
url = data.get("url")
if url:
from app.routes.user_sites import validate_webhook_url
is_safe, error_msg = validate_webhook_url(url)
if not is_safe:
return jsonify({"errors": [f"Invalid URL: {error_msg}"]}), 400
success = update_destination(dest_id, **data)
if not success:
return jsonify({"error": "No valid fields to update"}), 400
updated = get_destination(dest_id)
return jsonify(
{
"destination": {
"id": updated["id"],
"name": updated["name"],
"type": updated["type"],
"type_name": INTEGRATION_INFO.get(updated["type"], {}).get("name", updated["type"]),
"icon": INTEGRATION_INFO.get(updated["type"], {}).get("icon", "🔌"),
"url_display": _mask_url(updated.get("url")),
"enabled": updated.get("enabled", True),
}
}
)
@integrations_bp.route("/<int:dest_id>", methods=["DELETE"])
@login_required
def delete_destination_route(site_id, dest_id):
"""Delete a destination."""
site = _resolve_site(site_id)
dest = _resolve_destination(dest_id, site)
delete_destination(dest_id)
return jsonify({"deleted": True, "name": dest["name"]})
# ─── Available integrations ───────────────────────────────────────────────────
@integrations_bp.route("/available", methods=["GET"])
@login_required
def available_integrations(site_id):
"""List all available integration types with setup guides."""
_resolve_site(site_id)
result = []
for key, info in INTEGRATION_INFO.items():
guide = get_setup_guide(key)
result.append(
{
"type": key,
"name": info["name"],
"icon": info["icon"],
"description": info["description"],
"config_fields": guide.get("config_fields", []) if guide else [],
"requires_url": guide.get("requires_url", True) if guide else True,
"setup_guide": guide.get("steps", []) if guide else [],
}
)
return jsonify(result)
@integrations_bp.route("/<int:dest_id>/guide", methods=["GET"])
@login_required
def integration_guide(site_id, dest_id):
"""Get the setup guide for a destination's integration type."""
site = _resolve_site(site_id)
dest = _resolve_destination(dest_id, site)
guide = get_setup_guide(dest["type"])
info = INTEGRATION_INFO.get(dest["type"], {})
if not guide:
return jsonify({"error": "No setup guide available for this type"}), 404
return jsonify(
{
"type": dest["type"],
"name": info.get("name", dest["type"]),
"icon": info.get("icon", "🔌"),
"steps": guide.get("steps", []),
"config_fields": guide.get("config_fields", []),
}
)
# ─── Test delivery ────────────────────────────────────────────────────────────
@integrations_bp.route("/<int:dest_id>/test", methods=["POST"])
@login_required
def test_destination(site_id, dest_id):
"""Test a destination by firing it with the most recent submission."""
site = _resolve_site(site_id)
dest = _resolve_destination(dest_id, site)
# Find the most recent submission for this site
recent = list_submissions(limit=1, site_filter=site_id, password_hash=_get_site_owner_password_hash(site_id))
if not recent:
return jsonify(
{
"error": "No submissions found for this site — can't test without data",
"tip": "Submit a form entry first, then run this test again.",
}
), 400
submission = recent[0]
field_config = parse_site_fields(site)
try:
fire_destination(dest, site, submission, field_config)
except Exception as e:
return jsonify({"error": f"Test failed: {str(e)}"}), 500
# Re-fetch destination to get updated stats
updated = get_destination(dest_id)
return jsonify(
{
"success": True,
"message": f"Test payload sent to {dest['type']} destination",
"submission_id": submission["id"],
"last_status": updated.get("last_status"),
"last_delivered_at": updated.get("last_delivered_at"),
}
)
# ─── Export blueprint (separate prefix) ──────────────────────────────
export_bp = Blueprint(
"integrations_export",
__name__,
url_prefix="/sites/<int:site_id>/submissions",
)
# ─── CSV / JSON export ────────────────────────────────────────────────
@export_bp.route("/export")
@login_required
def export_submissions(site_id):
"""Export all submissions for a site as CSV or JSON."""
site = _resolve_site(site_id)
format_type = (request.args.get("format") or "json").lower().strip()
if format_type not in ("csv", "json"):
format_type = "json"
# Fetch all submissions for the site (unlimited)
all_submissions = list_submissions(
limit=10000, site_filter=site_id, password_hash=_get_site_owner_password_hash(site_id)
)
field_config = parse_site_fields(site)
if format_type == "csv":
return _export_csv(all_submissions, field_config, site)
else:
return _export_json(all_submissions, field_config)
def _build_row(submission, field_config):
"""Build a flat dict from a submission, merging hardcoded + dynamic fields."""
row = {
"id": submission.get("id", ""),
"submitted_at": submission.get("submitted_at", ""),
}
# Hardcoded field aliases
_aliases = {
"customer_name": "name",
"customer_phone": "phone",
"customer_email": "email",
"customer_equipment": "equipment",
"customer_message": "message",
}
for db_key, alias in _aliases.items():
val = submission.get(db_key)
if val:
row[alias] = val
# Merge dynamic data from JSON column
raw_data = submission.get("data")
if raw_data:
try:
dynamic = json.loads(raw_data)
if isinstance(dynamic, dict):
row.update(dynamic)
except (json.JSONDecodeError, TypeError):
pass
# Filter by field_config if provided
if field_config:
field_names = {f.get("name", f.get("key", "")) for f in field_config}
filtered = {k: v for k, v in row.items() if k in field_names}
# Always keep id and submitted_at
filtered["id"] = row["id"]
filtered["submitted_at"] = row["submitted_at"]
row = filtered
return row
def _export_csv(submissions, field_config, site):
"""Generate a CSV response from submissions."""
if not submissions:
csv_buf = io.StringIO()
writer = csv.writer(csv_buf)
writer.writerow(["No submissions found"])
csv_buf.seek(0)
return Response(
csv_buf.getvalue(),
mimetype="text/csv",
headers={"Content-Disposition": f"attachment; filename=submissions_{site['id']}.csv"},
)
# Determine headers: use field_config order, or just use keys from first submission
first_row = _build_row(submissions[0], field_config)
headers = list(first_row.keys())
csv_buf = io.StringIO()
writer = csv.writer(csv_buf)
writer.writerow(headers)
for sub in submissions:
row = _build_row(sub, field_config)
writer.writerow([row.get(h, "") for h in headers])
csv_buf.seek(0)
return Response(
csv_buf.getvalue(),
mimetype="text/csv",
headers={"Content-Disposition": f"attachment; filename=submissions_{site['id']}.csv"},
)
def _export_json(submissions, field_config):
"""Generate a JSON response from submissions."""
rows = [_build_row(sub, field_config) for sub in submissions]
return jsonify(
{
"count": len(rows),
"submissions": rows,
}
)