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
66 changes: 65 additions & 1 deletion DEPLOY.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,77 @@ Target domain: **localmate.crewcircle.co** (dashboard) / **api.localmate.crewcir
| Layer | Host | Domain | Trigger |
|---|---|---|---|
| Next.js dashboard | **Vercel** | `localmate.crewcircle.co` | GitHub Actions `deploy.yml` on push to `main` |
| FastAPI + APScheduler | **Docker Compose + Caddy** on shared DigitalOcean Sydney droplet (`170.64.183.45`) | `api.localmate.crewcircle.co` | `scripts/deploy_backend.sh` via Doppler |
| FastAPI (web) + arq worker + scheduler + Redis | **Docker Compose + Caddy** on shared DigitalOcean Sydney droplet (`170.64.183.45`) | `api.localmate.crewcircle.co` | `scripts/deploy_backend.sh` via Doppler |
| Postgres | **Supabase** | `*.supabase.co` | Pulumi provisioner |
| Secrets | **Doppler** `localmate/prd` (inherits `crewcircle-master/prod`) | — | Doppler CLI |
| DNS | **Cloudflare** | `crewcircle.co` zone | `scripts/finish_infra.sh` via Doppler |

The backend shares the same droplet as **TaxFlowAI** — Caddy routes `api.localmate.crewcircle.co` to the localmate backend container and `api.taxflow.crewcircle.com.au` to the taxflow backend container. No new droplet needed.

### Task queue containers (Phase 0)

The backend image runs in three roles, all from `deploy/docker-compose.yml`:

| Container | `WORKER_ROLE` | Command | Purpose |
|---|---|---|---|
| `localmate-backend` | `web` | `uvicorn main:app` | FastAPI API; enqueues jobs onto arq. Does **not** run APScheduler. |
| `localmate-worker` | `worker` | `arq worker.WorkerSettings` | Executes queued tasks (inbound webhook processing + durable outbound sends). Scale with `docker compose up -d --scale localmate-worker=N`. |
| `localmate-scheduler` | `scheduler` | `uvicorn main:app` | **Single-active** enqueue-only APScheduler (NOT HA). Each cron trigger pushes an arq job. Must remain a single instance so crons fire exactly once. |
| `localmate-redis` | — | `redis-server --appendonly yes --maxmemory 128mb --maxmemory-policy noeviction` | arq broker/result store. Persistent (AOF), non-evicting, memory-capped (128MB) on the shared 1GB droplet. Named volume `localmate-redis-data`. |

**Durability:** inbound webhooks (Stripe/GBP/menu) are persisted to `webhook_events` before enqueue; a reconciliation job (`reconcile_webhooks`, every 5 min) re-enqueues rows stuck `pending`. Exhausted retries (inbound and outbound) land in `dead_letter`. Redis is non-evicting so job state is never silently dropped.

### New Doppler vars (`localmate/prd`)

- `REDIS_URL` — `redis://localmate-redis:6379/0` (compose service DNS name). **Required.**
- `WORKER_ROLE` — set per-container in `docker-compose.yml` (`web`/`worker`/`scheduler`); no need to set in Doppler.
- `STRIPE_PORTAL_CONFIG_ID` — optional `bpc_...`; empty uses the Stripe dashboard default portal config. Used by the Phase 1 billing portal endpoint. Create the portal configuration in Stripe (enable subscription update to the localmate price, payment-method update, cancellation) — one-time dashboard/API step.
- `DASHBOARD_URL` — `https://localmate.crewcircle.co`; used as the Stripe portal `return_url` base.
- `SUPABASE_JWT_SECRET` — **Phase 1 (required for billing).** The Supabase project's JWT secret (Settings → API → JWT Settings in the Supabase dashboard). Used to verify dashboard bearer tokens and derive `client_id` from the authenticated user (tenant-auth binding, C8/D20). When unset, the legacy `require_auth` falls back to anonymous for old callers, but the strict `/billing/*` endpoints reject with 401 — so it must be set in prd.

### Migrations (Phase 0)

Apply `supabase/migrations/010_webhook_events.sql`, `011_dead_letter.sql`, `012_booking_credentials.sql` on top of `009`. Each new table has RLS enabled with a `service_role_all` policy (matches `001_clients.sql`). The sequence applies cleanly on top of `009` and is reversible (drop the two new tables / the added `clients` columns).

### Rollback (queue)

`docker compose stop localmate-worker localmate-scheduler localmate-redis` reverts to web-only. Inbound webhooks still persist to `webhook_events` as `pending` (nothing crashes); processing resumes when the worker + Redis come back and the reconciler re-enqueues the backlog.

### Phase 0 staged integration gate (C10 — MUST PASS before merge)

Unit tests mock Redis/Postgres, which cannot prove the durability behaviour. Before the Phase 0 PR merges to prod, run the executable gate — it spins up throwaway `postgres:15-alpine` + `redis:7-alpine` containers and **exits non-zero on any failed acceptance check** (nothing to eyeball):

```bash
bash scripts/phase0_staging_gate.sh
```

It asserts, and FAILS the build otherwise:

1. **Migrations** — Phase 0 migrations `010`/`011`/`012` apply cleanly on the reconstructed baseline; `webhook_events` + `dead_letter` exist with a `service_role_all` RLS policy; `012` added the five booking-credential columns on `clients`. (A pre-existing baseline-migration quirk already live in prod is a warning, not a Phase 0 failure.)
2. **Worker drain / restart** — an enqueued inbound job is drained by a real arq worker to `status='done'`.
3. **Retry + dead-letter** — a permanently-failing job is retried up to `MAX_TRIES` (real `arq.Retry` semantics) and then lands as exactly one `dead_letter` row.
4. **Reconciliation** — a stale `pending` row is re-enqueued and a stale `processing` lease is reset to `pending` and re-enqueued.
5. **Exactly-once / dedupe** — enqueuing the same event twice with the deterministic `_job_id` yields a single job (arq drops the duplicate).
6. **Redis durability (D9-A)** — the gate refuses to pass unless Redis is `noeviction` + AOF, so job state is never silently dropped.

Record the gate's `ALL PHASE 0 STAGING GATE CHECKS PASSED` output on the PR before merge.

### Phase 1 — billing visibility + tenant auth

**Migration:** apply `supabase/migrations/013_user_client_map.sql` on top of `012`. It creates `user_client_map (user_id text primary key, client_id uuid references clients(id))` with RLS + a `service_role_all` policy (matches `001_clients.sql`). One row per Supabase auth user → owned client; populate it at signup (or backfill existing users to their `clients` row by email). Until a user has a row, `/billing/*` returns 403.

**Tenant-auth binding (C8/D20):** `GET /billing/usage` and `POST /billing/portal` derive `client_id` from the authenticated identity via `user_client_map` (never from the request body/query), so a user can only reach their own tenant. `SUPABASE_JWT_SECRET` must be set in Doppler — with it empty these endpoints reject 401 (the legacy `require_auth` still falls back to anonymous for old non-client-scoped callers).

**One-time Stripe Billing Portal configuration:**
1. Stripe Dashboard → Settings → Billing → Customer portal → Add configuration (or via API: `stripe.billing_portal.Configuration.create(...)`).
2. Enable: subscription update (to the localmate price), payment-method update, subscription cancellation.
3. Copy the configuration id (`bpc_...`).
4. Set `STRIPE_PORTAL_CONFIG_ID` in Doppler (`localmate/prd`). Leave empty to use the Stripe dashboard default portal config.
5. `DASHBOARD_URL` controls where the user lands after leaving the portal (the `return_url`).

Portal-driven subscription changes (plan/card/cancel) flow back through the existing Stripe webhook handlers (`customer.subscription.updated` / `deleted`) which already update `clients.subscription_status` — no new reconciliation path.


---

## Phase A — You (one-time, ~10 min)
Expand Down
7 changes: 7 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@
SUPABASE_URL=
SUPABASE_ANON_KEY=
SUPABASE_SERVICE_ROLE_KEY=
# Phase 1 — tenant auth (Supabase JWT secret; required for /billing/* client-scoped endpoints)
SUPABASE_JWT_SECRET=
STRIPE_SECRET_KEY=
STRIPE_PRICE_ID=
STRIPE_WEBHOOK_SECRET=
Expand All @@ -17,6 +19,11 @@ DATAFORSEO_LOGIN=
DATAFORSEO_PASSWORD=
GBP_CLIENT_ID=
GBP_CLIENT_SECRET=
# Phase 0 — task queue / worker / billing portal
REDIS_URL=redis://localhost:6379/0
WORKER_ROLE=web
STRIPE_PORTAL_CONFIG_ID=
DASHBOARD_URL=http://localhost:3000
BASE_DOMAIN=crewcircle.com.au
PROJECT_ID=local-biz-au
ENVIRONMENT=prod
8 changes: 8 additions & 0 deletions backend/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,15 @@ COPY . .
# DATAFORSEO_LOGIN, DATAFORSEO_PASSWORD
# GBP_CLIENT_ID, GBP_CLIENT_SECRET
# BASE_DOMAIN, PROJECT_ID
# REDIS_URL (arq broker; e.g. redis://localmate-redis:6379/0)
# WORKER_ROLE (web|worker|scheduler; set per-container in compose)
# STRIPE_PORTAL_CONFIG_ID (optional; Stripe billing portal config bpc_...)
# DASHBOARD_URL (optional; portal return_url base)
# SENTRY_DSN (optional)
#
# Run modes (same image, different command):
# web / scheduler : uv run uvicorn main:app --host 0.0.0.0 --port 8000
# worker : uv run arq worker.WorkerSettings

EXPOSE 8000

Expand Down
6 changes: 6 additions & 0 deletions backend/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,12 @@ class Settings(BaseSettings):
square_environment: str = "sandbox"
supabase_jwt_secret: str = ""

# --- Phase 0: queue / worker / billing-portal infra ---
redis_url: str = "redis://localhost:6379/0"
worker_role: str = "web" # "web" | "worker" | "scheduler"
dashboard_url: str = "" # Stripe portal return_url base
stripe_portal_config_id: str = "" # Stripe portal configuration id (bpc_...)

class Config:
env_file = ".env.local"
env_file_encoding = "utf-8"
Expand Down
105 changes: 69 additions & 36 deletions backend/jobs/trial_emails.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,17 @@
"""Trial milestone emails — day1, day7, day13, and expired.

APScheduler-compatible async job. Runs hourly (wired via CronTrigger(hour=9)
in scheduler.py). Uses trial_emails_sent table for idempotency: one email
per (client_id, day_number) pair, ever.
APScheduler-compatible async job. Runs daily (wired via CronTrigger(hour=9) in
scheduler.py, which enqueues the ``run_trial_emails_daily`` arq task). Uses the
``trial_emails_sent`` table for idempotency: one email per (client_id,
day_number) pair, ever.

Outbound sends are durable (C4): instead of calling the Resend senders directly,
this job enqueues a ``send_email_task`` arq job (the Phase 0 durable wrapper) so
a Resend transport failure retries with backoff and dead-letters on exhaustion
rather than being silently swallowed. The idempotency row is recorded BEFORE the
enqueue and rolled back if the enqueue fails, so a failed enqueue can never leave
a row that would block a future send (which would cause a silent drop); a
permanently-failed send lands in ``dead_letter`` for operator replay.
"""

import logging
Expand All @@ -11,22 +20,16 @@
import pytz

from db import get_db
from services.resend_email import (
send_trial_day1_email,
send_trial_day7_email,
send_trial_day13_email,
send_trial_expired_email,
)

AEST = pytz.timezone("Australia/Sydney")
logger = logging.getLogger(__name__)

# days_since_trial_start -> (send_fn, needs_client_id_arg).
# send_trial_expired_email differs: (to, business_name) only — handled separately.
_DAY_DISPATCH: dict[int, tuple] = {
1: (send_trial_day1_email, True),
7: (send_trial_day7_email, True),
13: (send_trial_day13_email, True),
# days_since_trial_start -> Resend email kind (see services.resend_email._content).
# The expired email differs (no client_id arg) and is handled separately.
_DAY_KIND: dict[int, str] = {
1: "trial_day1",
7: "trial_day7",
13: "trial_day13",
}


Expand All @@ -48,7 +51,28 @@ def _record_send(db, client_id: str, day_number: int) -> None:
).execute()


async def run_trial_emails() -> None:
def _unrecord_send(db, client_id: str, day_number: int) -> None:
"""Remove a tentatively-recorded idempotency row (enqueue-failure rollback)."""
(
db.table("trial_emails_sent")
.delete()
.eq("client_id", client_id)
.eq("day_number", day_number)
.execute()
)


async def run_trial_emails(pool) -> None:
"""Enqueue durable trial-milestone emails for all active-trial clients.

``pool`` is the arq Redis pool (``ctx["redis"]`` from the
``run_trial_emails_daily`` cron task). When absent the run is skipped — the
next run with a live pool picks it up.
"""
if pool is None:
logger.warning("run_trial_emails: no arq pool available — skipping this run")
return

db = get_db()
now = datetime.now(AEST)

Expand Down Expand Up @@ -79,46 +103,55 @@ async def run_trial_emails() -> None:
).astimezone(AEST)
days_since = (now.date() - trial_start.date()).days

if days_since in _DAY_DISPATCH:
fn, _needs_id = _DAY_DISPATCH[days_since]
if days_since in _DAY_KIND:
if not _already_sent(db, client_id, days_since):
# Record the idempotency row FIRST so a crash between
# enqueue-success and record does not cause a duplicate
# delivery on the next run. If the enqueue fails, roll the
# row back so the next run can retry (no silent drop).
_record_send(db, client_id, days_since)
try:
await fn(client["email"], client["business_name"], client_id)
_record_send(db, client_id, days_since)
await pool.enqueue_job(
"send_email_task",
_DAY_KIND[days_since],
client["email"],
client["business_name"],
client_id,
)
sent += 1
logger.info(
"Sent day-%d email to %s (%s)",
"Enqueued day-%d email to %s (%s)",
days_since, client["email"], client_id,
)
except Exception as exc:
logger.error(
"send_trial_day%d_email failed for %s: %s",
days_since, client_id, exc,
)
except Exception:
_unrecord_send(db, client_id, days_since)
raise

raw_ends = client.get("trial_ends_at")
if raw_ends:
trial_ends = datetime.fromisoformat(
raw_ends.replace("Z", "+00:00")
).astimezone(AEST)
if now > trial_ends and not _already_sent(db, client_id, -1):
# Record-first / rollback-on-enqueue-failure (see above).
_record_send(db, client_id, -1)
try:
await send_trial_expired_email(
client["email"], client["business_name"]
await pool.enqueue_job(
"send_email_task",
"trial_expired",
client["email"],
client["business_name"],
)
_record_send(db, client_id, -1)
sent += 1
logger.info(
"Sent expired email to %s (%s)",
"Enqueued expired email to %s (%s)",
client["email"], client_id,
)
except Exception as exc:
logger.error(
"send_trial_expired_email failed for %s: %s",
client_id, exc,
)
except Exception:
_unrecord_send(db, client_id, -1)
raise

except Exception as exc:
logger.error("trial_emails: unexpected error for client %s: %s", client_id, exc)

logger.info("run_trial_emails: processed %d clients, sent %d emails", processed, sent)
logger.info("run_trial_emails: processed %d clients, enqueued %d emails", processed, sent)
38 changes: 33 additions & 5 deletions backend/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
from config import settings
from db import init_db
from scheduler import create_scheduler
from routers import auth, webhooks, drafts
from routers import auth, webhooks, drafts, billing

logging.basicConfig(level=logging.INFO)

Expand All @@ -23,11 +23,38 @@
@asynccontextmanager
async def lifespan(app: FastAPI):
await init_db()
scheduler = create_scheduler()
scheduler.start()
app.state.scheduler = scheduler
app.state.scheduler = None
app.state.arq = None

from task_queue import get_arq_pool

log = logging.getLogger(__name__)
if settings.worker_role == "scheduler":
# Dedicated single-active scheduler container: enqueue-only APScheduler
# + an arq pool to push jobs. NOT HA — must be a single instance (C4).
app.state.arq = await get_arq_pool()
scheduler = create_scheduler()
scheduler.start()
app.state.scheduler = scheduler
log.info("Started enqueue-only scheduler (role=scheduler)")
else:
# web role: create an arq pool for enqueuing from request handlers.
# Do NOT start APScheduler here (prevents duplicate cron fire across
# web replicas). The 'worker' role runs via the arq CLI, not uvicorn.
try:
app.state.arq = await get_arq_pool()
except Exception as e:
log.warning("arq pool init failed (web role): %s", e)

yield
scheduler.shutdown()

if app.state.scheduler is not None:
app.state.scheduler.shutdown()
if app.state.arq is not None:
try:
await app.state.arq.close()
except Exception:
pass


app = FastAPI(title="LocalMate", lifespan=lifespan)
Expand All @@ -48,6 +75,7 @@ async def lifespan(app: FastAPI):
app.include_router(auth.router, prefix="/auth")
app.include_router(webhooks.router, prefix="/webhooks")
app.include_router(drafts.router, prefix="/drafts")
app.include_router(billing.router, prefix="/billing")
try:
from routers import approve
app.include_router(approve.router, prefix="/approve")
Expand Down
Loading
Loading