-
{invitation.email}
+
+ {invitation.email}
+ {invitation.role === "viewer" && (
+ read-only
+ )}
+
{expired
? "Expired"
@@ -708,6 +785,375 @@ function Field({
);
}
+function DataSourcesPanel({
+ organizationId,
+ dataSources,
+}: {
+ organizationId: string;
+ dataSources: DashboardDataSource[];
+}) {
+ const router = useRouter();
+ const [message, setMessage] = useState(null);
+ const [pending, startTransition] = useTransition();
+
+ function syncAll() {
+ setMessage(null);
+ startTransition(async () => {
+ const result = await syncAllOrganizationAudience({ organizationId });
+ if (!result.ok) {
+ setMessage(result.error);
+ return;
+ }
+ setMessage(
+ `Synced ${result.imported} emails (${result.added} new)` +
+ (result.failed ? `, ${result.failed} source(s) failed` : ""),
+ );
+ router.refresh();
+ });
+ }
+
+ return (
+
+
+ User data sources
+
+
+
+ Connect each project's database (Supabase or Turso). Syncing pulls
+ every user email into this org's deduplicated audience.
+
+ {dataSources.length === 0 ? (
+
No data sources yet.
+ ) : (
+
+ {dataSources.map((source) => (
+
+
+
+
{source.label}
+
+ {source.kind}
+ {!source.enabled && off }
+
+
+ {source.last_sync_error ? (
+ {source.last_sync_error}
+ ) : source.last_synced_at ? (
+ <>
+ {source.last_sync_count ?? 0} emails ·{" "}
+ {new Date(source.last_synced_at).toLocaleString()}
+ >
+ ) : (
+ "Never synced"
+ )}
+
+
+
+
+
+ ))}
+
+ )}
+
+
+ {pending ? "Syncing..." : "Sync all"}
+
+ {message && (
+ {message}
+ )}
+
+
+
+
+ );
+}
+
+function DataSourceActions({
+ organizationId,
+ sourceId,
+}: {
+ organizationId: string;
+ sourceId: string;
+}) {
+ const router = useRouter();
+ const [message, setMessage] = useState(null);
+ const [pending, startTransition] = useTransition();
+
+ function sync() {
+ setMessage(null);
+ startTransition(async () => {
+ const result = await syncOrganizationDataSource({ organizationId, sourceId });
+ setMessage(result.ok ? `+${result.added} new` : result.error);
+ if (result.ok) router.refresh();
+ });
+ }
+
+ function remove() {
+ setMessage(null);
+ startTransition(async () => {
+ const result = await deleteOrganizationDataSource({ organizationId, sourceId });
+ if (!result.ok) {
+ setMessage(result.error);
+ return;
+ }
+ router.refresh();
+ });
+ }
+
+ return (
+
+
+ {pending ? "..." : "Sync now"}
+
+
+ Delete
+
+ {message &&
{message}
}
+
+ );
+}
+
+type DataSourceKind = "supabase" | "turso";
+
+function DataSourceForm({ organizationId }: { organizationId: string }) {
+ const router = useRouter();
+ const [kind, setKind] = useState("supabase");
+ const [mode, setMode] = useState<"auth_users" | "table">("auth_users");
+ const [message, setMessage] = useState(null);
+ const [pending, startTransition] = useTransition();
+
+ function submit(event: FormEvent) {
+ event.preventDefault();
+ const formEl = event.currentTarget;
+ const form = new FormData(formEl);
+ setMessage(null);
+ startTransition(async () => {
+ const result = await saveOrganizationDataSource({
+ organizationId,
+ label: String(form.get("label") ?? ""),
+ kind,
+ supabaseUrl: String(form.get("supabaseUrl") ?? ""),
+ serviceRoleKey: String(form.get("serviceRoleKey") ?? ""),
+ sourceMode: mode,
+ tableName: String(form.get("tableName") ?? ""),
+ emailColumn: String(form.get("emailColumn") ?? ""),
+ tursoUrl: String(form.get("tursoUrl") ?? ""),
+ authToken: String(form.get("authToken") ?? ""),
+ emailQuery: String(form.get("emailQuery") ?? ""),
+ });
+ setMessage(result.ok ? "Source saved." : result.error);
+ if (result.ok) {
+ formEl.reset();
+ router.refresh();
+ }
+ });
+ }
+
+ return (
+
+ );
+}
+
+function AudiencePanel({
+ organizationId,
+ audienceStats,
+}: {
+ organizationId: string;
+ audienceStats: DashboardAudienceStats;
+}) {
+ const [message, setMessage] = useState(null);
+ const [pending, startTransition] = useTransition();
+ const active = audienceStats.total - audienceStats.unsubscribed;
+
+ function send(previewTo?: string) {
+ return (event: FormEvent) => {
+ event.preventDefault();
+ const form = new FormData(event.currentTarget);
+ const subject = String(form.get("subject") ?? "");
+ const html = String(form.get("html") ?? "");
+ setMessage(null);
+ startTransition(async () => {
+ const result = await sendOrganizationAudienceBlast({
+ organizationId,
+ subject,
+ html,
+ previewTo,
+ });
+ if (!result.ok) {
+ setMessage(result.error);
+ return;
+ }
+ setMessage(
+ previewTo
+ ? `Preview sent (${result.sent} sent, ${result.skipped} skipped).`
+ : `Done — ${result.sent} sent, ${result.failed} failed, ${result.skipped} skipped of ${result.total}.`,
+ );
+ });
+ };
+ }
+
+ return (
+
+
+ Mass email ({active} contacts)
+
+
+
+ Sends to {active} active contacts via this org's default email sender.
+ {audienceStats.unsubscribed > 0 &&
+ ` ${audienceStats.unsubscribed} unsubscribed are excluded.`}{" "}
+ Every message includes a one-click unsubscribe link.
+
+
+ {message &&
{message}
}
+
+
+ );
+}
+
+function CampaignForm({
+ onSubmitSend,
+ pending,
+}: {
+ onSubmitSend: (previewTo?: string) => (event: FormEvent) => void;
+ pending: boolean;
+}) {
+ const [previewTo, setPreviewTo] = useState("");
+ const formRef = useRef(null);
+
+ return (
+
+ );
+}
+
function DangerZonePanel({
orgId,
orgs,
diff --git a/app/(app)/dashboard/page.tsx b/app/(app)/dashboard/page.tsx
index d88237dc..a48e789b 100644
--- a/app/(app)/dashboard/page.tsx
+++ b/app/(app)/dashboard/page.tsx
@@ -9,6 +9,8 @@ import {
OrgDashboardControls,
ProjectOrgMoveControl,
type DashboardSenderConfig,
+ type DashboardDataSource,
+ type DashboardAudienceStats,
type DashboardOrg,
type DashboardOrgTeam,
} from "./org-controls";
@@ -92,7 +94,16 @@ export default async function DashboardPage({
? projectsQuery.eq("organization_id", selectedOrgId).or(accessFilter)
: projectsQuery.or(accessFilter);
- const [{ data: projectsRaw }, { data: audits }, counts, senderConfigs, orgTeam] = await Promise.all([
+ const ownerOrgId = selectedOrg?.role === "owner" ? selectedOrgId : null;
+ const [
+ { data: projectsRaw },
+ { data: audits },
+ counts,
+ senderConfigs,
+ dataSources,
+ audienceStats,
+ orgTeam,
+ ] = await Promise.all([
scopedProjectsQuery,
supabase
.from("audits")
@@ -101,7 +112,9 @@ export default async function DashboardPage({
.order("created_at", { ascending: false })
.limit(10),
countByStatus(supabase, user!.id, accessFilter, selectedOrgId),
- fetchSenderConfigs(supabase, selectedOrg?.role === "owner" ? selectedOrgId : null),
+ fetchSenderConfigs(supabase, ownerOrgId),
+ fetchDataSources(supabase, ownerOrgId),
+ fetchAudienceStats(supabase, ownerOrgId),
fetchOrgTeam(selectedOrg && isOrgWideRole(selectedOrg.role) ? selectedOrgId : null),
]);
const projects = ((projectsRaw ?? []) as unknown) as DashboardProject[];
@@ -141,6 +154,8 @@ export default async function DashboardPage({
orgs={orgs as DashboardOrg[]}
selectedOrgId={selectedOrgId}
senderConfigs={senderConfigs}
+ dataSources={dataSources}
+ audienceStats={audienceStats}
orgTeam={orgTeam}
/>
)}
@@ -519,3 +534,43 @@ async function fetchSenderConfigs(
}
return ((data ?? []) as unknown) as DashboardSenderConfig[];
}
+
+async function fetchDataSources(
+ supabase: Awaited>,
+ organizationId: string | null,
+): Promise {
+ if (!organizationId) return [];
+ const { data, error } = await supabase
+ .from("organization_data_sources")
+ .select("id,label,kind,enabled,last_synced_at,last_sync_count,last_sync_error")
+ .eq("organization_id", organizationId)
+ .order("created_at", { ascending: false });
+ if (error) {
+ if (missingOrgSchema(error)) return [];
+ throw error;
+ }
+ return ((data ?? []) as unknown) as DashboardDataSource[];
+}
+
+async function fetchAudienceStats(
+ supabase: Awaited>,
+ organizationId: string | null,
+): Promise {
+ if (!organizationId) return { total: 0, unsubscribed: 0 };
+ const [{ count: total, error: totalErr }, { count: unsub, error: unsubErr }] =
+ await Promise.all([
+ supabase
+ .from("organization_audience_contacts")
+ .select("id", { count: "exact", head: true })
+ .eq("organization_id", organizationId),
+ supabase
+ .from("organization_audience_contacts")
+ .select("id", { count: "exact", head: true })
+ .eq("organization_id", organizationId)
+ .not("unsubscribed_at", "is", null),
+ ]);
+ if (totalErr || unsubErr) {
+ if (missingOrgSchema(totalErr ?? unsubErr)) return { total: 0, unsubscribed: 0 };
+ }
+ return { total: total ?? 0, unsubscribed: unsub ?? 0 };
+}
diff --git a/app/actions/org-members.ts b/app/actions/org-members.ts
index 59599ad5..91e16a53 100644
--- a/app/actions/org-members.ts
+++ b/app/actions/org-members.ts
@@ -34,14 +34,20 @@ async function requireOrgOwner(orgId: string) {
return { ok: true as const, user, org, supabase, svc };
}
+export type OrgInviteRole = "member" | "viewer";
+
export async function inviteOrgMember(
orgId: string,
email: string,
+ role: OrgInviteRole = "member",
): Promise<{ ok: boolean; error?: string }> {
const normalized = email.trim().toLowerCase();
if (!normalized.includes("@")) {
return { ok: false, error: "Invalid email address." };
}
+ if (role !== "member" && role !== "viewer") {
+ return { ok: false, error: "Invalid role." };
+ }
const ctx = await requireOrgOwner(orgId);
if (!ctx.ok) return { ok: false, error: ctx.error };
@@ -73,7 +79,7 @@ export async function inviteOrgMember(
if ((existingMember as { role: OrgRole }).role === "project_member") {
const { error } = await svc
.from("organization_members")
- .update({ role: "member" })
+ .update({ role })
.eq("id", (existingMember as { id: string }).id);
if (error) return { ok: false, error: error.message };
revalidatePath("/dashboard");
@@ -92,7 +98,7 @@ export async function inviteOrgMember(
const { data: invitation, error: insertErr } = await svc
.from("organization_invitations")
- .insert({ organization_id: orgId, email: normalized, invited_by: user.id })
+ .insert({ organization_id: orgId, email: normalized, invited_by: user.id, role })
.select("token")
.single();
@@ -169,13 +175,13 @@ export async function removeOrgMember(
.update({ role: "project_member" })
.eq("organization_id", orgId)
.eq("user_id", userId)
- .eq("role", "member")
+ .in("role", ["member", "viewer"])
: await ctx.svc
.from("organization_members")
.delete()
.eq("organization_id", orgId)
.eq("user_id", userId)
- .eq("role", "member");
+ .in("role", ["member", "viewer"]);
if (error) return { ok: false, error: error.message };
@@ -191,6 +197,31 @@ export async function removeOrgMember(
return { ok: true };
}
+// Change an existing org member between full member and read-only viewer.
+// Owner-only. Never touches 'owner' or 'project_member' rows.
+export async function setOrgMemberRole(
+ orgId: string,
+ userId: string,
+ role: OrgInviteRole,
+): Promise<{ ok: boolean; error?: string }> {
+ if (role !== "member" && role !== "viewer") {
+ return { ok: false, error: "Invalid role." };
+ }
+ const ctx = await requireOrgOwner(orgId);
+ if (!ctx.ok) return { ok: false, error: ctx.error };
+
+ const { error } = await ctx.svc
+ .from("organization_members")
+ .update({ role })
+ .eq("organization_id", orgId)
+ .eq("user_id", userId)
+ .in("role", ["member", "viewer"]);
+ if (error) return { ok: false, error: error.message };
+
+ revalidatePath("/dashboard");
+ return { ok: true };
+}
+
export type OrgTeamMember = {
id: string;
user_id: string;
@@ -202,6 +233,7 @@ export type OrgTeamMember = {
export type OrgPendingInvitation = {
id: string;
email: string;
+ role: OrgInviteRole;
expires_at: string;
created_at: string;
};
@@ -264,7 +296,7 @@ export async function listOrgTeam(orgId: string): Promise<
const { data: invitationsRaw } = isOwner
? await svc
.from("organization_invitations")
- .select("id, email, expires_at, created_at")
+ .select("id, email, role, expires_at, created_at")
.eq("organization_id", orgId)
.is("accepted_at", null)
.order("created_at", { ascending: false })
diff --git a/app/actions/orgs.ts b/app/actions/orgs.ts
index 2fbbbd1a..2ee50357 100644
--- a/app/actions/orgs.ts
+++ b/app/actions/orgs.ts
@@ -5,6 +5,9 @@ import { createClient } from "@/lib/supabase/server";
import { serviceClient } from "@/lib/supabase/service";
import { getOrCreateDefaultOrg } from "@/lib/orgs";
import { encryptSecret } from "@/lib/sp/vault";
+import { syncDataSource, syncAllForOrg } from "@/lib/audience/sync";
+import { sendOrgAudienceBlast } from "@/lib/audience/blast";
+import { assertReadOnlySelect } from "@/lib/audience/connectors";
type Ok = { ok: true } & (T extends undefined ? {} : T);
type Err = { ok: false; error: string };
@@ -301,6 +304,197 @@ export async function deleteOrganizationOutreachConfig(input: {
return { ok: true };
}
+// --- Org audience: connected project databases + mass email ----------------
+
+async function requireOrgOwner(
+ organizationId: string,
+): Promise<{ ok: true; userId: string } | Err> {
+ const supabase = await createClient();
+ const {
+ data: { user },
+ } = await supabase.auth.getUser();
+ if (!user) return { ok: false, error: "Not authenticated." };
+ const svc = serviceClient();
+ const { data: member } = await svc
+ .from("organization_members")
+ .select("id")
+ .eq("organization_id", organizationId)
+ .eq("user_id", user.id)
+ .eq("role", "owner")
+ .maybeSingle();
+ if (!member) return { ok: false, error: "You must own this org." };
+ return { ok: true, userId: user.id };
+}
+
+export async function saveOrganizationDataSource(input: {
+ organizationId: string;
+ label: string;
+ kind: "supabase" | "turso";
+ // Supabase
+ supabaseUrl?: string;
+ serviceRoleKey?: string;
+ sourceMode?: "auth_users" | "table";
+ tableName?: string;
+ emailColumn?: string;
+ // Turso
+ tursoUrl?: string;
+ authToken?: string;
+ emailQuery?: string;
+}): Promise | Err> {
+ const owner = await requireOrgOwner(input.organizationId);
+ if (!owner.ok) return owner;
+
+ const label = input.label.trim().replace(/\s+/g, " ").slice(0, 80);
+ if (!label) return { ok: false, error: "Label is required." };
+
+ const patch: Record = {
+ organization_id: input.organizationId,
+ created_by: owner.userId,
+ label,
+ kind: input.kind,
+ enabled: true,
+ supabase_url: null,
+ enc_service_role_key: null,
+ source_mode: null,
+ table_name: null,
+ email_column: null,
+ turso_url: null,
+ enc_auth_token: null,
+ email_query: null,
+ };
+
+ try {
+ if (input.kind === "supabase") {
+ const url = clean(input.supabaseUrl);
+ if (!url) return { ok: false, error: "Supabase URL is required." };
+ const mode = input.sourceMode === "table" ? "table" : "auth_users";
+ patch.supabase_url = url;
+ patch.source_mode = mode;
+ if (mode === "table") {
+ const table = clean(input.tableName);
+ const column = clean(input.emailColumn);
+ if (!table || !column) {
+ return { ok: false, error: "Table name and email column are required for table mode." };
+ }
+ patch.table_name = table;
+ patch.email_column = column;
+ }
+ const key = clean(input.serviceRoleKey);
+ if (key) patch.enc_service_role_key = encryptSecret(key);
+ else return { ok: false, error: "Service role key is required." };
+ } else {
+ const url = clean(input.tursoUrl);
+ const query = clean(input.emailQuery);
+ if (!url) return { ok: false, error: "Turso URL is required." };
+ if (!query) return { ok: false, error: "Email query is required." };
+ const guard = assertReadOnlySelect(query);
+ if (guard) return { ok: false, error: guard };
+ patch.turso_url = url;
+ patch.email_query = query;
+ const token = clean(input.authToken);
+ if (token) patch.enc_auth_token = encryptSecret(token);
+ }
+ } catch (error) {
+ return {
+ ok: false,
+ error: error instanceof Error ? error.message : "Could not encrypt source credentials.",
+ };
+ }
+
+ const svc = serviceClient();
+ const { data, error } = await svc
+ .from("organization_data_sources")
+ .insert(patch)
+ .select("id")
+ .single();
+ if (error) return { ok: false, error: error.message };
+
+ revalidatePath("/dashboard");
+ return { ok: true, id: data.id as string };
+}
+
+export async function deleteOrganizationDataSource(input: {
+ organizationId: string;
+ sourceId: string;
+}): Promise {
+ const owner = await requireOrgOwner(input.organizationId);
+ if (!owner.ok) return owner;
+
+ const svc = serviceClient();
+ const { error } = await svc
+ .from("organization_data_sources")
+ .delete()
+ .eq("id", input.sourceId)
+ .eq("organization_id", input.organizationId);
+ if (error) return { ok: false, error: error.message };
+
+ revalidatePath("/dashboard");
+ return { ok: true };
+}
+
+export async function syncOrganizationDataSource(input: {
+ organizationId: string;
+ sourceId: string;
+}): Promise | Err> {
+ const owner = await requireOrgOwner(input.organizationId);
+ if (!owner.ok) return owner;
+
+ const result = await syncDataSource(input.organizationId, input.sourceId);
+ if (!result.ok) return { ok: false, error: result.error ?? "Sync failed." };
+
+ revalidatePath("/dashboard");
+ return { ok: true, imported: result.imported, added: result.added };
+}
+
+export async function syncAllOrganizationAudience(input: {
+ organizationId: string;
+}): Promise | Err> {
+ const owner = await requireOrgOwner(input.organizationId);
+ if (!owner.ok) return owner;
+
+ const results = await syncAllForOrg(input.organizationId);
+ const imported = results.reduce((n, r) => n + r.imported, 0);
+ const added = results.reduce((n, r) => n + r.added, 0);
+ const failed = results.filter((r) => !r.ok).length;
+
+ revalidatePath("/dashboard");
+ return { ok: true, imported, added, failed };
+}
+
+export async function sendOrganizationAudienceBlast(input: {
+ organizationId: string;
+ subject: string;
+ html: string;
+ previewTo?: string;
+}): Promise | Err> {
+ const owner = await requireOrgOwner(input.organizationId);
+ if (!owner.ok) return owner;
+
+ const subject = input.subject.trim();
+ const html = input.html.trim();
+ if (!subject) return { ok: false, error: "Subject is required." };
+ if (!html) return { ok: false, error: "Message body is required." };
+
+ const preview = clean(input.previewTo);
+ const result = await sendOrgAudienceBlast({
+ organizationId: input.organizationId,
+ subject,
+ html,
+ createdBy: owner.userId,
+ previewTo: preview ?? undefined,
+ });
+ if (!result.ok) return { ok: false, error: result.error ?? "Send failed." };
+
+ revalidatePath("/dashboard");
+ return {
+ ok: true,
+ total: result.total,
+ sent: result.sent,
+ failed: result.failed,
+ skipped: result.skipped,
+ };
+}
+
export async function deleteOrganization(input: {
orgId: string;
}): Promise {
diff --git a/app/invite/org/[token]/page.tsx b/app/invite/org/[token]/page.tsx
index 3a12c36a..fc37e5b2 100644
--- a/app/invite/org/[token]/page.tsx
+++ b/app/invite/org/[token]/page.tsx
@@ -15,7 +15,7 @@ export default async function OrgInvitePage({
const { data: inv } = await svc
.from("organization_invitations")
- .select("id, organization_id, email, expires_at, accepted_at")
+ .select("id, organization_id, email, role, expires_at, accepted_at")
.eq("token", token)
.maybeSingle();
@@ -105,16 +105,19 @@ export default async function OrgInvitePage({
.eq("user_id", user.id)
.maybeSingle();
+ const invitedRole =
+ (inv as { role?: string }).role === "viewer" ? "viewer" : "member";
+
if (!existing) {
await svc.from("organization_members").insert({
organization_id: inv.organization_id,
user_id: user.id,
- role: "member",
+ role: invitedRole,
});
} else if ((existing as { role?: string }).role === "project_member") {
await svc
.from("organization_members")
- .update({ role: "member" })
+ .update({ role: invitedRole })
.eq("id", (existing as { id: string }).id);
}
diff --git a/app/unsubscribe/org/[token]/page.tsx b/app/unsubscribe/org/[token]/page.tsx
new file mode 100644
index 00000000..8dac4984
--- /dev/null
+++ b/app/unsubscribe/org/[token]/page.tsx
@@ -0,0 +1,42 @@
+import Link from "next/link";
+import { unsubscribeOrgAudienceByToken } from "@/lib/marketing";
+
+export const metadata = { title: "Unsubscribe" };
+export const dynamic = "force-dynamic";
+
+export default async function OrgUnsubscribePage({
+ params,
+}: {
+ params: Promise<{ token: string }>;
+}) {
+ const { token } = await params;
+ const result = await unsubscribeOrgAudienceByToken(token);
+
+ return (
+
+
+ {result.ok ? "You're unsubscribed" : "Unsubscribe link not recognized"}
+
+ {result.ok ? (
+
+ {result.email ? (
+ <>
+ {result.email} won't receive any more emails
+ from us.
+ >
+ ) : (
+ <>You won't receive any more emails from us.>
+ )}
+
+ ) : (
+
+ We couldn't find a subscription for that link. It may have
+ already been unsubscribed, or the link may be malformed.
+
+ )}
+
+ ← Back to CrawlProof
+
+
+ );
+}
diff --git a/lib/audience/blast.ts b/lib/audience/blast.ts
new file mode 100644
index 00000000..422cc068
--- /dev/null
+++ b/lib/audience/blast.ts
@@ -0,0 +1,186 @@
+import "server-only";
+import { serviceClient } from "@/lib/supabase/service";
+import { env } from "@/lib/env";
+import { sendOutreachEmail, type OutreachConfig } from "@/lib/outreach";
+
+export type BlastResult = {
+ ok: boolean;
+ total: number;
+ sent: number;
+ failed: number;
+ skipped: number;
+ error?: string;
+ campaignId?: string;
+};
+
+// Mass-email an org's deduped audience through its configured email sender
+// (SMTP or Resend — whichever is the org's default email config). Every
+// message carries a one-click unsubscribe footer + List-Unsubscribe headers,
+// and any globally-unsubscribed address (marketing_contacts) is skipped.
+export async function sendOrgAudienceBlast(input: {
+ organizationId: string;
+ subject: string;
+ html: string;
+ createdBy?: string | null;
+ // Send only to this address (must be a real, active audience contact).
+ // Used to preview a campaign before going live.
+ previewTo?: string;
+ perSecond?: number;
+}): Promise {
+ const svc = serviceClient();
+ const limitPerSec = Math.max(1, Math.min(input.perSecond ?? 5, 20));
+ const delayMs = Math.ceil(1000 / limitPerSec);
+
+ // 1. Resolve the org's default email sender config.
+ const { data: configRow, error: configErr } = await svc
+ .from("organization_outreach_configs")
+ .select(
+ "id,provider,from_email,reply_to,smtp_host,smtp_port,smtp_secure,enc_smtp_user,enc_smtp_pass,enc_api_key",
+ )
+ .eq("organization_id", input.organizationId)
+ .eq("channel", "email")
+ .eq("enabled", true)
+ .eq("is_default", true)
+ .maybeSingle();
+ if (configErr) return zero("Could not load sender config: " + configErr.message);
+ if (!configRow) {
+ return zero("No default email sender configured. Add an SMTP or Resend sender first.");
+ }
+ const config = configRow as unknown as OutreachConfig;
+
+ // 2. Build the suppression set: globally-unsubscribed marketing contacts.
+ const suppressed = await loadSuppressed(svc);
+
+ // 3. Page through active audience contacts.
+ const contacts = await loadActiveContacts(svc, input.organizationId, input.previewTo);
+ if (contacts.length === 0) {
+ return { ok: true, total: 0, sent: 0, failed: 0, skipped: 0 };
+ }
+
+ let sent = 0;
+ let failed = 0;
+ let skipped = 0;
+
+ for (const c of contacts) {
+ if (suppressed.has(c.email)) {
+ skipped += 1;
+ continue;
+ }
+ const unsubUrl = `${env.siteUrl}/unsubscribe/org/${c.unsubscribe_token}`;
+ const html = `${input.html}${unsubscribeFooter(unsubUrl)}`;
+ const res = await sendOutreachEmail({
+ to: c.email,
+ subject: input.subject,
+ body: htmlToText(input.html) + `\n\nUnsubscribe: ${unsubUrl}`,
+ html,
+ headers: {
+ "List-Unsubscribe": `<${unsubUrl}>`,
+ "List-Unsubscribe-Post": "List-Unsubscribe=One-Click",
+ },
+ config,
+ });
+ if (res.sent) sent += 1;
+ else {
+ failed += 1;
+ console.warn("[audience.blast] send failed", { to: c.email, error: res.error });
+ }
+ if (delayMs > 0) await new Promise((r) => setTimeout(r, delayMs));
+ }
+
+ // 4. Record the campaign (skip the audit row for previews).
+ let campaignId: string | undefined;
+ if (!input.previewTo) {
+ const { data: campaign } = await svc
+ .from("organization_email_campaigns")
+ .insert({
+ organization_id: input.organizationId,
+ created_by: input.createdBy ?? null,
+ sender_config_id: configRow.id,
+ subject: input.subject,
+ sent_count: sent,
+ failed_count: failed,
+ skipped_count: skipped,
+ })
+ .select("id")
+ .maybeSingle();
+ campaignId = campaign?.id as string | undefined;
+ }
+
+ return { ok: true, total: contacts.length, sent, failed, skipped, campaignId };
+}
+
+function zero(error: string): BlastResult {
+ return { ok: false, total: 0, sent: 0, failed: 0, skipped: 0, error };
+}
+
+async function loadActiveContacts(
+ svc: ReturnType,
+ organizationId: string,
+ previewTo?: string,
+): Promise> {
+ if (previewTo) {
+ const email = previewTo.trim().toLowerCase();
+ const { data } = await svc
+ .from("organization_audience_contacts")
+ .select("email,unsubscribe_token,unsubscribed_at")
+ .eq("organization_id", organizationId)
+ .ilike("email", email)
+ .maybeSingle();
+ if (!data || data.unsubscribed_at) return [];
+ return [{ email: data.email as string, unsubscribe_token: data.unsubscribe_token as string }];
+ }
+
+ const out: Array<{ email: string; unsubscribe_token: string }> = [];
+ const pageSize = 1000;
+ for (let from = 0; ; from += pageSize) {
+ const { data, error } = await svc
+ .from("organization_audience_contacts")
+ .select("email,unsubscribe_token")
+ .eq("organization_id", organizationId)
+ .is("unsubscribed_at", null)
+ .range(from, from + pageSize - 1);
+ if (error || !data || data.length === 0) break;
+ for (const r of data) {
+ out.push({ email: r.email as string, unsubscribe_token: r.unsubscribe_token as string });
+ }
+ if (data.length < pageSize) break;
+ }
+ return out;
+}
+
+async function loadSuppressed(svc: ReturnType): Promise> {
+ const set = new Set();
+ const pageSize = 1000;
+ for (let from = 0; ; from += pageSize) {
+ const { data, error } = await svc
+ .from("marketing_contacts")
+ .select("email")
+ .not("unsubscribed_at", "is", null)
+ .range(from, from + pageSize - 1);
+ if (error || !data || data.length === 0) break;
+ for (const r of data) set.add(String(r.email).trim().toLowerCase());
+ if (data.length < pageSize) break;
+ }
+ return set;
+}
+
+function unsubscribeFooter(unsubUrl: string): string {
+ return `
+
+ You're receiving this because you have an account on one of our products.
+ Unsubscribe .
+
`;
+}
+
+function htmlToText(html: string): string {
+ return html
+ .replace(/<\s*br\s*\/?>/gi, "\n")
+ .replace(/<\/(p|div|h[1-6]|li)>/gi, "\n")
+ .replace(/<[^>]+>/g, "")
+ .replace(/ /g, " ")
+ .replace(/&/g, "&")
+ .replace(/</g, "<")
+ .replace(/>/g, ">")
+ .replace(/\n{3,}/g, "\n\n")
+ .trim();
+}
diff --git a/lib/audience/connectors.ts b/lib/audience/connectors.ts
new file mode 100644
index 00000000..790d2166
--- /dev/null
+++ b/lib/audience/connectors.ts
@@ -0,0 +1,141 @@
+import { createClient as createSb } from "@supabase/supabase-js";
+import { createClient as createTurso } from "@libsql/client";
+
+// Connectors pull every user email out of a project's own backing database.
+// Supabase and Turso are the two hosted stores used across the CrawlProof
+// org's projects. Each connector returns a flat, normalized email list; the
+// caller (lib/audience/sync.ts) is responsible for dedup + persistence.
+
+export type DataSourceRow = {
+ id: string;
+ organization_id: string;
+ kind: "supabase" | "turso";
+ // Supabase
+ supabase_url: string | null;
+ source_mode: "auth_users" | "table" | null;
+ table_name: string | null;
+ email_column: string | null;
+ // Turso
+ turso_url: string | null;
+ email_query: string | null;
+ // Decrypted at call time by the caller.
+ serviceRoleKey?: string | null;
+ authToken?: string | null;
+};
+
+export type FetchResult = { emails: string[]; error?: string };
+
+// Normalize a raw value to a deliverable email or null. Lowercase + trim,
+// require a single "@" with something either side, drop anything obviously
+// junk. Mirrors lib/marketing.ts#normalize, slightly stricter.
+export function normalizeEmail(value: unknown): string | null {
+ if (typeof value !== "string") return null;
+ const e = value.trim().toLowerCase();
+ if (!e || e.length > 254) return null;
+ const at = e.indexOf("@");
+ if (at <= 0 || at !== e.lastIndexOf("@") || at === e.length - 1) return null;
+ if (/\s/.test(e)) return null;
+ if (!e.slice(at + 1).includes(".")) return null;
+ return e;
+}
+
+function dedupeNormalize(values: unknown[]): string[] {
+ const seen = new Set();
+ for (const v of values) {
+ const e = normalizeEmail(v);
+ if (e) seen.add(e);
+ }
+ return [...seen];
+}
+
+// Guard a user-supplied Turso query: must be a single read-only SELECT.
+// The DB belongs to the org owner, so this is defense-in-depth — it stops a
+// fat-fingered destructive statement, not a determined attacker.
+export function assertReadOnlySelect(query: string): string | null {
+ const q = query.trim().replace(/;+\s*$/, ""); // allow one trailing semicolon
+ if (!q) return "Query is empty.";
+ if (q.includes(";")) return "Query must be a single statement (no semicolons).";
+ if (!/^\s*(with|select)\b/i.test(q)) return "Query must be a SELECT.";
+ if (/\b(insert|update|delete|drop|alter|create|attach|detach|pragma|replace|truncate|vacuum|reindex)\b/i.test(q)) {
+ return "Query may only read data (no write/DDL keywords).";
+ }
+ return null;
+}
+
+export async function fetchSupabaseEmails(src: DataSourceRow): Promise {
+ if (!src.supabase_url) return { emails: [], error: "Supabase URL is required." };
+ if (!src.serviceRoleKey) return { emails: [], error: "Service role key is required." };
+
+ const sb = createSb(src.supabase_url, src.serviceRoleKey, {
+ auth: { autoRefreshToken: false, persistSession: false },
+ });
+
+ if (src.source_mode === "table") {
+ if (!src.table_name || !src.email_column) {
+ return { emails: [], error: "Table name and email column are required for table mode." };
+ }
+ const collected: unknown[] = [];
+ const pageSize = 1000;
+ for (let from = 0; ; from += pageSize) {
+ const { data, error } = await sb
+ .from(src.table_name)
+ .select(src.email_column)
+ .range(from, from + pageSize - 1);
+ if (error) return { emails: [], error: error.message };
+ if (!data || data.length === 0) break;
+ for (const row of data) {
+ collected.push((row as unknown as Record)[src.email_column]);
+ }
+ if (data.length < pageSize) break;
+ }
+ return { emails: dedupeNormalize(collected) };
+ }
+
+ // Default: read auth.users via the admin API — uniform across every
+ // Supabase project regardless of its public schema.
+ const collected: unknown[] = [];
+ const perPage = 1000;
+ for (let page = 1; ; page += 1) {
+ const { data, error } = await sb.auth.admin.listUsers({ page, perPage });
+ if (error) return { emails: [], error: error.message };
+ const users = data?.users ?? [];
+ if (users.length === 0) break;
+ for (const u of users) collected.push(u.email);
+ if (users.length < perPage) break;
+ }
+ return { emails: dedupeNormalize(collected) };
+}
+
+export async function fetchTursoEmails(src: DataSourceRow): Promise {
+ if (!src.turso_url) return { emails: [], error: "Turso URL is required." };
+ if (!src.email_query) return { emails: [], error: "Email query is required." };
+ const guardError = assertReadOnlySelect(src.email_query);
+ if (guardError) return { emails: [], error: guardError };
+
+ const client = createTurso({
+ url: src.turso_url,
+ authToken: src.authToken ?? undefined,
+ });
+ try {
+ const result = await client.execute(src.email_query);
+ // Prefer an "email" column if present, else the first column.
+ const cols = result.columns ?? [];
+ const emailIdx = cols.findIndex((c) => c?.toLowerCase() === "email");
+ const values: unknown[] = result.rows.map((row) => {
+ if (emailIdx >= 0) return (row as unknown as unknown[])[emailIdx];
+ const arr = row as unknown as unknown[];
+ return arr[0];
+ });
+ return { emails: dedupeNormalize(values) };
+ } catch (error) {
+ return { emails: [], error: error instanceof Error ? error.message : "Turso query failed." };
+ } finally {
+ client.close();
+ }
+}
+
+export async function fetchEmailsForSource(src: DataSourceRow): Promise {
+ if (src.kind === "supabase") return fetchSupabaseEmails(src);
+ if (src.kind === "turso") return fetchTursoEmails(src);
+ return { emails: [], error: `Unknown data source kind: ${src.kind}` };
+}
diff --git a/lib/audience/sync.ts b/lib/audience/sync.ts
new file mode 100644
index 00000000..0bb1ecba
--- /dev/null
+++ b/lib/audience/sync.ts
@@ -0,0 +1,145 @@
+import "server-only";
+import { serviceClient } from "@/lib/supabase/service";
+import { decryptSecret } from "@/lib/sp/vault";
+import { fetchEmailsForSource, type DataSourceRow } from "./connectors";
+
+export type SyncResult = {
+ ok: boolean;
+ sourceId: string;
+ imported: number; // distinct emails returned by the connector
+ added: number; // rows newly inserted into the audience
+ error?: string;
+};
+
+// Pull emails from one connected data source and upsert them into the org's
+// deduped audience. Existing rows (incl. their unsubscribe state) are
+// preserved — re-syncing never resurrects an unsubscribed contact.
+export async function syncDataSource(
+ organizationId: string,
+ sourceId: string,
+): Promise {
+ const svc = serviceClient();
+ const { data: row, error } = await svc
+ .from("organization_data_sources")
+ .select(
+ "id,organization_id,kind,enabled,supabase_url,enc_service_role_key,source_mode,table_name,email_column,turso_url,enc_auth_token,email_query",
+ )
+ .eq("id", sourceId)
+ .eq("organization_id", organizationId)
+ .maybeSingle();
+
+ if (error || !row) {
+ return { ok: false, sourceId, imported: 0, added: 0, error: error?.message ?? "Source not found." };
+ }
+
+ const src: DataSourceRow = {
+ id: row.id,
+ organization_id: row.organization_id,
+ kind: row.kind,
+ supabase_url: row.supabase_url,
+ source_mode: row.source_mode,
+ table_name: row.table_name,
+ email_column: row.email_column,
+ turso_url: row.turso_url,
+ email_query: row.email_query,
+ };
+
+ try {
+ if (row.enc_service_role_key) src.serviceRoleKey = decryptSecret(row.enc_service_role_key);
+ if (row.enc_auth_token) src.authToken = decryptSecret(row.enc_auth_token);
+ } catch {
+ const msg = "Could not decrypt source credentials (check SOCIAL_VAULT_KEY).";
+ await recordSyncError(svc, sourceId, msg);
+ return { ok: false, sourceId, imported: 0, added: 0, error: msg };
+ }
+
+ const fetched = await fetchEmailsForSource(src);
+ if (fetched.error) {
+ await recordSyncError(svc, sourceId, fetched.error);
+ return { ok: false, sourceId, imported: 0, added: 0, error: fetched.error };
+ }
+
+ let added = 0;
+ if (fetched.emails.length > 0) {
+ // Which of these already exist for this org? Page through to avoid an
+ // unbounded IN list, then insert only the new ones. The unique index on
+ // (organization_id, lower(email)) is the final safety net under races.
+ const existing = await loadExistingEmails(svc, organizationId);
+ const fresh = fetched.emails.filter((e) => !existing.has(e));
+ for (let i = 0; i < fresh.length; i += 500) {
+ const chunk = fresh.slice(i, i + 500).map((email) => ({
+ organization_id: organizationId,
+ source_id: sourceId,
+ email,
+ }));
+ const { error: insErr, count } = await svc
+ .from("organization_audience_contacts")
+ .upsert(chunk, { onConflict: "organization_id,email", ignoreDuplicates: true, count: "exact" });
+ if (insErr) {
+ // Fall back to per-row insert ignoring conflicts if the bulk upsert
+ // conflict target doesn't match the functional index.
+ for (const r of chunk) {
+ const { error: rowErr } = await svc.from("organization_audience_contacts").insert(r);
+ if (!rowErr) added += 1;
+ }
+ continue;
+ }
+ added += count ?? chunk.length;
+ }
+ }
+
+ await svc
+ .from("organization_data_sources")
+ .update({
+ last_synced_at: new Date().toISOString(),
+ last_sync_count: fetched.emails.length,
+ last_sync_error: null,
+ })
+ .eq("id", sourceId);
+
+ return { ok: true, sourceId, imported: fetched.emails.length, added };
+}
+
+export async function syncAllForOrg(organizationId: string): Promise {
+ const svc = serviceClient();
+ const { data } = await svc
+ .from("organization_data_sources")
+ .select("id")
+ .eq("organization_id", organizationId)
+ .eq("enabled", true);
+ const results: SyncResult[] = [];
+ for (const s of data ?? []) {
+ results.push(await syncDataSource(organizationId, s.id as string));
+ }
+ return results;
+}
+
+async function loadExistingEmails(
+ svc: ReturnType,
+ organizationId: string,
+): Promise> {
+ const emails = new Set();
+ const pageSize = 1000;
+ for (let from = 0; ; from += pageSize) {
+ const { data, error } = await svc
+ .from("organization_audience_contacts")
+ .select("email")
+ .eq("organization_id", organizationId)
+ .range(from, from + pageSize - 1);
+ if (error || !data || data.length === 0) break;
+ for (const r of data) emails.add(String(r.email).toLowerCase());
+ if (data.length < pageSize) break;
+ }
+ return emails;
+}
+
+async function recordSyncError(
+ svc: ReturnType,
+ sourceId: string,
+ message: string,
+): Promise {
+ await svc
+ .from("organization_data_sources")
+ .update({ last_synced_at: new Date().toISOString(), last_sync_error: message.slice(0, 500) })
+ .eq("id", sourceId);
+}
diff --git a/lib/marketing.ts b/lib/marketing.ts
index 5c0492c0..eafd4edc 100644
--- a/lib/marketing.ts
+++ b/lib/marketing.ts
@@ -103,3 +103,20 @@ export async function unsubscribeByToken(
if (error || !data) return { ok: false };
return { ok: true, email: data.email as string };
}
+
+// Unsubscribe an imported org-audience contact (mass-email recipient). Marks
+// the per-org row suppressed; the contact won't receive future org blasts.
+export async function unsubscribeOrgAudienceByToken(
+ token: string,
+): Promise<{ ok: boolean; email?: string }> {
+ if (!token || token.length < 8) return { ok: false };
+ const svc = serviceClient();
+ const { data, error } = await svc
+ .from("organization_audience_contacts")
+ .update({ unsubscribed_at: new Date().toISOString() })
+ .eq("unsubscribe_token", token)
+ .select("email")
+ .maybeSingle();
+ if (error || !data) return { ok: false };
+ return { ok: true, email: data.email as string };
+}
diff --git a/lib/orgs.ts b/lib/orgs.ts
index b2e5acb3..1e14e97c 100644
--- a/lib/orgs.ts
+++ b/lib/orgs.ts
@@ -8,10 +8,12 @@ export type OrgSummary = {
role: OrgRole;
};
-export type OrgRole = "owner" | "member" | "project_member";
+export type OrgRole = "owner" | "member" | "viewer" | "project_member";
-export function isOrgWideRole(role: OrgRole | null | undefined): role is "owner" | "member" {
- return role === "owner" || role === "member";
+export function isOrgWideRole(
+ role: OrgRole | null | undefined,
+): role is "owner" | "member" | "viewer" {
+ return role === "owner" || role === "member" || role === "viewer";
}
export function missingOrgSchema(error: unknown) {
diff --git a/lib/outreach.ts b/lib/outreach.ts
index 04c0b970..d9444b03 100644
--- a/lib/outreach.ts
+++ b/lib/outreach.ts
@@ -38,6 +38,12 @@ export async function sendOutreachEmail(input: {
to: string;
subject: string;
body: string;
+ // Optional pre-rendered HTML. When omitted, HTML is derived from `body`.
+ // Used by audience campaigns to send real HTML + unsubscribe footer.
+ html?: string;
+ // Extra mail headers (e.g. List-Unsubscribe). Passed to both SMTP and
+ // Resend so native unsubscribe buttons render.
+ headers?: Record;
replyTo?: string | null;
config?: OutreachConfig | null;
}): Promise {
@@ -59,10 +65,15 @@ export async function sendOutreachEmail(input: {
};
}
+ // Honor an explicit per-org provider choice: a "resend" config must never
+ // be hijacked by a global SMTP_HOST env. With no config we keep the legacy
+ // default of falling back to the global SMTP host.
const smtpHost =
- input.config?.provider === "smtp" && input.config.smtp_host
- ? input.config.smtp_host
- : env.smtpHost;
+ input.config?.provider === "resend"
+ ? ""
+ : input.config?.provider === "smtp"
+ ? input.config.smtp_host ?? env.smtpHost
+ : env.smtpHost;
if (smtpHost) {
try {
const transporter = nodemailer.createTransport({
@@ -88,7 +99,8 @@ export async function sendOutreachEmail(input: {
replyTo: input.config?.reply_to ?? input.replyTo ?? undefined,
subject: input.subject,
text: input.body,
- html: paragraphHtml(input.body),
+ html: input.html ?? paragraphHtml(input.body),
+ headers: input.headers,
});
return {
sent: true,
@@ -116,7 +128,8 @@ export async function sendOutreachEmail(input: {
replyTo: input.config?.reply_to ?? input.replyTo ?? undefined,
subject: input.subject,
text: input.body,
- html: paragraphHtml(input.body),
+ html: input.html ?? paragraphHtml(input.body),
+ headers: input.headers,
});
if (result.error) {
return { sent: false, provider: "resend", error: String(result.error) };
diff --git a/package-lock.json b/package-lock.json
index a9f6b03d..9cf813cd 100644
--- a/package-lock.json
+++ b/package-lock.json
@@ -10,6 +10,7 @@
"dependencies": {
"@anthropic-ai/sdk": "^0.95.2",
"@ip-location-db/geolite2-city-mmdb": "^2.3.2026052019",
+ "@libsql/client": "^0.17.3",
"@profullstack/autoblog": "github:profullstack/autoblog#75e54af",
"@profullstack/emailer": "^1.0.1",
"@profullstack/referrals": "^0.1.0",
@@ -1088,6 +1089,165 @@
"@jridgewell/sourcemap-codec": "^1.4.14"
}
},
+ "node_modules/@libsql/client": {
+ "version": "0.17.3",
+ "resolved": "https://registry.npmjs.org/@libsql/client/-/client-0.17.3.tgz",
+ "integrity": "sha512-HXk9wiAoJbKFbyBH4O+aEhN6ir5ERXuXvwE5OD2eR4/5RUa3Pw/8L9zrnVdU+iNJitRvisPWaIwmhkO3bH7giA==",
+ "license": "MIT",
+ "dependencies": {
+ "@libsql/core": "^0.17.3",
+ "@libsql/hrana-client": "^0.10.0",
+ "js-base64": "^3.7.5",
+ "libsql": "^0.5.28",
+ "promise-limit": "^2.7.0"
+ }
+ },
+ "node_modules/@libsql/core": {
+ "version": "0.17.3",
+ "resolved": "https://registry.npmjs.org/@libsql/core/-/core-0.17.3.tgz",
+ "integrity": "sha512-2UjK1i7JBkMduJo4WdvvBxMMvVJ31pArBZNONyz/GCJJAH+1UHat2X6vn10S/WpY5fKzIT98WqYFl2vzWRLOfg==",
+ "license": "MIT",
+ "dependencies": {
+ "js-base64": "^3.7.5"
+ }
+ },
+ "node_modules/@libsql/darwin-arm64": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/darwin-arm64/-/darwin-arm64-0.5.29.tgz",
+ "integrity": "sha512-K+2RIB1OGFPYQbfay48GakLhqf3ArcbHqPFu7EZiaUcRgFcdw8RoltsMyvbj5ix2fY0HV3Q3Ioa/ByvQdaSM0A==",
+ "cpu": [
+ "arm64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "darwin"
+ ]
+ },
+ "node_modules/@libsql/darwin-x64": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/darwin-x64/-/darwin-x64-0.5.29.tgz",
+ "integrity": "sha512-OtT+KFHsKFy1R5FVadr8FJ2Bb1mghtXTyJkxv0trocq7NuHntSki1eUbxpO5ezJesDvBlqFjnWaYYY516QNLhQ==",
+ "cpu": [
+ "x64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "darwin"
+ ]
+ },
+ "node_modules/@libsql/hrana-client": {
+ "version": "0.10.0",
+ "resolved": "https://registry.npmjs.org/@libsql/hrana-client/-/hrana-client-0.10.0.tgz",
+ "integrity": "sha512-OoA4EMqRAC7kn7V2P6EQqRcpZf2W+AjsNIyCizBg339Tq/aMC7sRnzs3SklderhmQWAqEzvv8A2vhxVmWpkVvw==",
+ "license": "MIT",
+ "dependencies": {
+ "@libsql/isomorphic-ws": "^0.1.5",
+ "js-base64": "^3.7.5"
+ }
+ },
+ "node_modules/@libsql/isomorphic-ws": {
+ "version": "0.1.5",
+ "resolved": "https://registry.npmjs.org/@libsql/isomorphic-ws/-/isomorphic-ws-0.1.5.tgz",
+ "integrity": "sha512-DtLWIH29onUYR00i0GlQ3UdcTRC6EP4u9w/h9LxpUZJWRMARk6dQwZ6Jkd+QdwVpuAOrdxt18v0K2uIYR3fwFg==",
+ "license": "MIT",
+ "dependencies": {
+ "@types/ws": "^8.5.4",
+ "ws": "^8.13.0"
+ }
+ },
+ "node_modules/@libsql/linux-arm-gnueabihf": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-arm-gnueabihf/-/linux-arm-gnueabihf-0.5.29.tgz",
+ "integrity": "sha512-CD4n4zj7SJTHso4nf5cuMoWoMSS7asn5hHygsDuhRl8jjjCTT3yE+xdUvI4J7zsyb53VO5ISh4cwwOtf6k2UhQ==",
+ "cpu": [
+ "arm"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/linux-arm-musleabihf": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-arm-musleabihf/-/linux-arm-musleabihf-0.5.29.tgz",
+ "integrity": "sha512-2Z9qBVpEJV7OeflzIR3+l5yAd4uTOLxklScYTwpZnkm2vDSGlC1PRlueLaufc4EFITkLKXK2MWBpexuNJfMVcg==",
+ "cpu": [
+ "arm"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/linux-arm64-gnu": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-arm64-gnu/-/linux-arm64-gnu-0.5.29.tgz",
+ "integrity": "sha512-gURBqaiXIGGwFNEaUj8Ldk7Hps4STtG+31aEidCk5evMMdtsdfL3HPCpvys+ZF/tkOs2MWlRWoSq7SOuCE9k3w==",
+ "cpu": [
+ "arm64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/linux-arm64-musl": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-arm64-musl/-/linux-arm64-musl-0.5.29.tgz",
+ "integrity": "sha512-fwgYZ0H8mUkyVqXZHF3mT/92iIh1N94Owi/f66cPVNsk9BdGKq5gVpoKO+7UxaNzuEH1roJp2QEwsCZMvBLpqg==",
+ "cpu": [
+ "arm64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/linux-x64-gnu": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-x64-gnu/-/linux-x64-gnu-0.5.29.tgz",
+ "integrity": "sha512-y14V0vY0nmMC6G0pHeJcEarcnGU2H6cm21ZceRkacWHvQAEhAG0latQkCtoS2njFOXiYIg+JYPfAoWKbi82rkg==",
+ "cpu": [
+ "x64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/linux-x64-musl": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/linux-x64-musl/-/linux-x64-musl-0.5.29.tgz",
+ "integrity": "sha512-gquqwA/39tH4pFl+J9n3SOMSymjX+6kZ3kWgY3b94nXFTwac9bnFNMffIomgvlFaC4ArVqMnOZD3nuJ3H3VO1w==",
+ "cpu": [
+ "x64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "linux"
+ ]
+ },
+ "node_modules/@libsql/win32-x64-msvc": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/@libsql/win32-x64-msvc/-/win32-x64-msvc-0.5.29.tgz",
+ "integrity": "sha512-4/0CvEdhi6+KjMxMaVbFM2n2Z44escBRoEYpR+gZg64DdetzGnYm8mcNLcoySaDJZNaBd6wz5DNdgRmcI4hXcg==",
+ "cpu": [
+ "x64"
+ ],
+ "license": "MIT",
+ "optional": true,
+ "os": [
+ "win32"
+ ]
+ },
"node_modules/@napi-rs/wasm-runtime": {
"version": "1.1.4",
"resolved": "https://registry.npmjs.org/@napi-rs/wasm-runtime/-/wasm-runtime-1.1.4.tgz",
@@ -1107,6 +1267,12 @@
"@emnapi/runtime": "^1.7.1"
}
},
+ "node_modules/@neon-rs/load": {
+ "version": "0.0.4",
+ "resolved": "https://registry.npmjs.org/@neon-rs/load/-/load-0.0.4.tgz",
+ "integrity": "sha512-kTPhdZyTQxB+2wpiRcFWrDcejc4JI6tkPuS7UZCG4l6Zvc5kU/gGQ/ozvHTh1XR5tS+UlfAfGuPajjzQjCiHCw==",
+ "license": "MIT"
+ },
"node_modules/@next/env": {
"version": "16.2.6",
"resolved": "https://registry.npmjs.org/@next/env/-/env-16.2.6.tgz",
@@ -2257,6 +2423,15 @@
"integrity": "sha512-zFDAD+tlpf2r4asuHEj0XH6pY6i0g5NeAHPn+15wk3BV6JA69eERFXC1gyGThDkVa1zCyKr5jox1+2LbV/AMLg==",
"license": "MIT"
},
+ "node_modules/@types/ws": {
+ "version": "8.18.1",
+ "resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz",
+ "integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==",
+ "license": "MIT",
+ "dependencies": {
+ "@types/node": "*"
+ }
+ },
"node_modules/@vitest/expect": {
"version": "4.1.6",
"resolved": "https://registry.npmjs.org/@vitest/expect/-/expect-4.1.6.tgz",
@@ -3473,6 +3648,12 @@
"jiti": "lib/jiti-cli.mjs"
}
},
+ "node_modules/js-base64": {
+ "version": "3.7.8",
+ "resolved": "https://registry.npmjs.org/js-base64/-/js-base64-3.7.8.tgz",
+ "integrity": "sha512-hNngCeKxIUQiEUN3GPJOkz4wF/YvdUdbNL9hsBcMQTkKzboD7T/q3OYOuuPZLUE6dBxSGpwhk5mwuDud7JVAow==",
+ "license": "BSD-3-Clause"
+ },
"node_modules/js-tokens": {
"version": "4.0.0",
"resolved": "https://registry.npmjs.org/js-tokens/-/js-tokens-4.0.0.tgz",
@@ -3513,6 +3694,47 @@
"url": "https://ko-fi.com/killymxi"
}
},
+ "node_modules/libsql": {
+ "version": "0.5.29",
+ "resolved": "https://registry.npmjs.org/libsql/-/libsql-0.5.29.tgz",
+ "integrity": "sha512-8lMP8iMgiBzzoNbAPQ59qdVcj6UaE/Vnm+fiwX4doX4Narook0a4GPKWBEv+CR8a1OwbfkgL18uBfBjWdF0Fzg==",
+ "cpu": [
+ "x64",
+ "arm64",
+ "wasm32",
+ "arm"
+ ],
+ "license": "MIT",
+ "os": [
+ "darwin",
+ "linux",
+ "win32"
+ ],
+ "dependencies": {
+ "@neon-rs/load": "^0.0.4",
+ "detect-libc": "2.0.2"
+ },
+ "optionalDependencies": {
+ "@libsql/darwin-arm64": "0.5.29",
+ "@libsql/darwin-x64": "0.5.29",
+ "@libsql/linux-arm-gnueabihf": "0.5.29",
+ "@libsql/linux-arm-musleabihf": "0.5.29",
+ "@libsql/linux-arm64-gnu": "0.5.29",
+ "@libsql/linux-arm64-musl": "0.5.29",
+ "@libsql/linux-x64-gnu": "0.5.29",
+ "@libsql/linux-x64-musl": "0.5.29",
+ "@libsql/win32-x64-msvc": "0.5.29"
+ }
+ },
+ "node_modules/libsql/node_modules/detect-libc": {
+ "version": "2.0.2",
+ "resolved": "https://registry.npmjs.org/detect-libc/-/detect-libc-2.0.2.tgz",
+ "integrity": "sha512-UX6sGumvvqSaXgdKGUsgZWqcUyIXZ/vZTrlRT/iobiKhGL0zL4d3osHj3uqllWJK+i+sixDS/3COVEOFbupFyw==",
+ "license": "Apache-2.0",
+ "engines": {
+ "node": ">=8"
+ }
+ },
"node_modules/lightningcss": {
"version": "1.32.0",
"resolved": "https://registry.npmjs.org/lightningcss/-/lightningcss-1.32.0.tgz",
@@ -4377,6 +4599,12 @@
"url": "https://github.com/prettier/prettier?sponsor=1"
}
},
+ "node_modules/promise-limit": {
+ "version": "2.7.0",
+ "resolved": "https://registry.npmjs.org/promise-limit/-/promise-limit-2.7.0.tgz",
+ "integrity": "sha512-7nJ6v5lnJsXwGprnGXga4wx6d1POjvi5Qmf1ivTRxTjH4Z/9Czja/UCMLVmB9N93GeWOU93XaFaEt6jbuoagNw==",
+ "license": "ISC"
+ },
"node_modules/prop-types": {
"version": "15.8.1",
"resolved": "https://registry.npmjs.org/prop-types/-/prop-types-15.8.1.tgz",
@@ -5409,6 +5637,27 @@
"node": ">=8"
}
},
+ "node_modules/ws": {
+ "version": "8.21.0",
+ "resolved": "https://registry.npmjs.org/ws/-/ws-8.21.0.tgz",
+ "integrity": "sha512-Vsp28b7DRcimFQvrqu2Wek3z1iYxDCWqHYB8Qsnk/S4RfaCQzPGPyBNuVjJV3cd6UiKtUtp6sNM77gWvzcCH+g==",
+ "license": "MIT",
+ "engines": {
+ "node": ">=10.0.0"
+ },
+ "peerDependencies": {
+ "bufferutil": "^4.0.1",
+ "utf-8-validate": ">=5.0.2"
+ },
+ "peerDependenciesMeta": {
+ "bufferutil": {
+ "optional": true
+ },
+ "utf-8-validate": {
+ "optional": true
+ }
+ }
+ },
"node_modules/zod": {
"version": "3.25.76",
"resolved": "https://registry.npmjs.org/zod/-/zod-3.25.76.tgz",
diff --git a/package.json b/package.json
index 474d98c1..d9a284ab 100644
--- a/package.json
+++ b/package.json
@@ -22,6 +22,7 @@
"dependencies": {
"@anthropic-ai/sdk": "^0.95.2",
"@ip-location-db/geolite2-city-mmdb": "^2.3.2026052019",
+ "@libsql/client": "^0.17.3",
"@profullstack/autoblog": "github:profullstack/autoblog#75e54af",
"@profullstack/emailer": "^1.0.1",
"@profullstack/referrals": "^0.1.0",
diff --git a/supabase/migrations/20260611140000_org_audience.sql b/supabase/migrations/20260611140000_org_audience.sql
new file mode 100644
index 00000000..ce618159
--- /dev/null
+++ b/supabase/migrations/20260611140000_org_audience.sql
@@ -0,0 +1,115 @@
+-- Org-wide audience: connect each project's backing database (Supabase or
+-- Turso), pull every user email into one deduped per-org list, and mass-email
+-- that list through the org's existing sender config (SMTP or Resend).
+--
+-- Mirrors the patterns in 20260606133000_prospects_outreach_configs.sql:
+-- RLS via public.is_org_owner, lx_set_updated_at trigger, and secrets held in
+-- enc_* columns (encrypted by lib/sp/vault.ts; plaintext secret columns stay
+-- null for new writes).
+
+-- 1. Connected project databases ------------------------------------------
+
+create table if not exists public.organization_data_sources (
+ id uuid primary key default gen_random_uuid(),
+ organization_id uuid not null references public.organizations(id) on delete cascade,
+ created_by uuid references public.profiles(id) on delete set null,
+ label text not null,
+ kind text not null check (kind in ('supabase', 'turso')),
+ enabled boolean not null default true,
+ -- Supabase
+ supabase_url text,
+ enc_service_role_key text,
+ source_mode text check (source_mode in ('auth_users', 'table')),
+ table_name text,
+ email_column text,
+ -- Turso (libSQL)
+ turso_url text,
+ enc_auth_token text,
+ email_query text,
+ -- Sync bookkeeping
+ last_synced_at timestamptz,
+ last_sync_count int,
+ last_sync_error text,
+ created_at timestamptz not null default now(),
+ updated_at timestamptz not null default now()
+);
+
+create index if not exists organization_data_sources_org_idx
+ on public.organization_data_sources(organization_id, enabled);
+
+alter table public.organization_data_sources enable row level security;
+
+drop policy if exists "organization_data_sources owner all"
+ on public.organization_data_sources;
+create policy "organization_data_sources owner all"
+ on public.organization_data_sources for all
+ using ((select public.is_org_owner(organization_id, auth.uid())))
+ with check ((select public.is_org_owner(organization_id, auth.uid())));
+
+drop trigger if exists organization_data_sources_set_updated_at
+ on public.organization_data_sources;
+create trigger organization_data_sources_set_updated_at
+ before update on public.organization_data_sources
+ for each row execute function public.lx_set_updated_at();
+
+-- 2. Deduped imported audience --------------------------------------------
+
+create table if not exists public.organization_audience_contacts (
+ id uuid primary key default gen_random_uuid(),
+ organization_id uuid not null references public.organizations(id) on delete cascade,
+ source_id uuid references public.organization_data_sources(id) on delete set null,
+ email text not null,
+ unsubscribe_token text not null unique default encode(gen_random_bytes(16), 'hex'),
+ unsubscribed_at timestamptz,
+ first_seen_at timestamptz not null default now(),
+ updated_at timestamptz not null default now()
+);
+
+-- Case-insensitive uniqueness per org so re-syncs upsert instead of dupe.
+create unique index if not exists organization_audience_contacts_org_email_idx
+ on public.organization_audience_contacts(organization_id, lower(email));
+
+create index if not exists organization_audience_contacts_active_idx
+ on public.organization_audience_contacts(organization_id)
+ where unsubscribed_at is null;
+
+alter table public.organization_audience_contacts enable row level security;
+
+drop policy if exists "organization_audience_contacts owner all"
+ on public.organization_audience_contacts;
+create policy "organization_audience_contacts owner all"
+ on public.organization_audience_contacts for all
+ using ((select public.is_org_owner(organization_id, auth.uid())))
+ with check ((select public.is_org_owner(organization_id, auth.uid())));
+
+drop trigger if exists organization_audience_contacts_set_updated_at
+ on public.organization_audience_contacts;
+create trigger organization_audience_contacts_set_updated_at
+ before update on public.organization_audience_contacts
+ for each row execute function public.lx_set_updated_at();
+
+-- 3. Campaign audit log ----------------------------------------------------
+
+create table if not exists public.organization_email_campaigns (
+ id uuid primary key default gen_random_uuid(),
+ organization_id uuid not null references public.organizations(id) on delete cascade,
+ created_by uuid references public.profiles(id) on delete set null,
+ sender_config_id uuid references public.organization_outreach_configs(id) on delete set null,
+ subject text not null,
+ sent_count int not null default 0,
+ failed_count int not null default 0,
+ skipped_count int not null default 0,
+ created_at timestamptz not null default now()
+);
+
+create index if not exists organization_email_campaigns_org_idx
+ on public.organization_email_campaigns(organization_id, created_at desc);
+
+alter table public.organization_email_campaigns enable row level security;
+
+drop policy if exists "organization_email_campaigns owner all"
+ on public.organization_email_campaigns;
+create policy "organization_email_campaigns owner all"
+ on public.organization_email_campaigns for all
+ using ((select public.is_org_owner(organization_id, auth.uid())))
+ with check ((select public.is_org_owner(organization_id, auth.uid())));
diff --git a/supabase/migrations/20260611150000_org_member_viewer_role.sql b/supabase/migrations/20260611150000_org_member_viewer_role.sql
new file mode 100644
index 00000000..0caada3e
--- /dev/null
+++ b/supabase/migrations/20260611150000_org_member_viewer_role.sql
@@ -0,0 +1,72 @@
+-- Org-level read-only members ("viewers").
+--
+-- Parallels the project-level viewer role (20260611120000_project_member_viewer_role.sql):
+-- owner/member = full org-wide access (read + write)
+-- viewer = org-wide READ-ONLY: sees every project in the org and its
+-- data, but cannot mutate any project or org settings
+-- project_member = org-visible marker only (access comes from project_members)
+--
+-- Reads are granted by widening the org-wide READ helpers (is_org_wide_member,
+-- is_project_member's org branch) to include 'viewer'. Writes are unchanged:
+-- is_project_editor and is_org_owner keep their ('owner','member') / owner
+-- checks, so viewers stay write-blocked everywhere.
+
+alter table public.organization_members
+ drop constraint if exists organization_members_role_check;
+alter table public.organization_members
+ add constraint organization_members_role_check
+ check (role in ('owner', 'member', 'viewer', 'project_member'));
+
+-- Carry the chosen role through the invitation so it is applied on accept.
+alter table public.organization_invitations
+ add column if not exists role text not null default 'member';
+alter table public.organization_invitations
+ drop constraint if exists organization_invitations_role_check;
+alter table public.organization_invitations
+ add constraint organization_invitations_role_check
+ check (role in ('member', 'viewer'));
+
+-- Org-wide READ access now includes viewers. Every usage of this helper is a
+-- SELECT policy (projects / integrations / event_outbox / webhook_events), so
+-- this grants reads only.
+create or replace function public.is_org_wide_member(p_org_id uuid, p_user_id uuid)
+returns boolean
+language sql
+security definer
+stable
+set search_path = public
+as $$
+ select exists(
+ select 1
+ from public.organization_members
+ where organization_id = p_org_id
+ and user_id = p_user_id
+ and role in ('owner', 'member', 'viewer')
+ )
+$$;
+
+-- Project READ via org membership now includes viewers. Project WRITES use
+-- is_project_editor (unchanged), which still excludes viewers.
+create or replace function public.is_project_member(p_project_id uuid, p_user_id uuid)
+returns boolean
+language sql
+security definer
+stable
+set search_path = public
+as $$
+ select exists(
+ select 1
+ from public.project_members pm
+ where pm.project_id = p_project_id
+ and pm.user_id = p_user_id
+ )
+ or exists(
+ select 1
+ from public.projects p
+ join public.organization_members om
+ on om.organization_id = p.organization_id
+ where p.id = p_project_id
+ and om.user_id = p_user_id
+ and om.role in ('owner', 'member', 'viewer')
+ )
+$$;
diff --git a/tests/audience-connectors.test.ts b/tests/audience-connectors.test.ts
new file mode 100644
index 00000000..7c41b184
--- /dev/null
+++ b/tests/audience-connectors.test.ts
@@ -0,0 +1,62 @@
+import { describe, it, expect } from "vitest";
+import { normalizeEmail, assertReadOnlySelect } from "@/lib/audience/connectors";
+
+describe("normalizeEmail", () => {
+ it("lowercases and trims", () => {
+ expect(normalizeEmail(" User@Example.COM ")).toBe("user@example.com");
+ });
+
+ it("rejects non-strings and blanks", () => {
+ expect(normalizeEmail(null)).toBeNull();
+ expect(normalizeEmail(undefined)).toBeNull();
+ expect(normalizeEmail(42)).toBeNull();
+ expect(normalizeEmail("")).toBeNull();
+ expect(normalizeEmail(" ")).toBeNull();
+ });
+
+ it("requires a single @ with a dotted domain", () => {
+ expect(normalizeEmail("nope")).toBeNull();
+ expect(normalizeEmail("a@b")).toBeNull(); // no dot in domain
+ expect(normalizeEmail("a@@b.com")).toBeNull();
+ expect(normalizeEmail("@example.com")).toBeNull();
+ expect(normalizeEmail("user@")).toBeNull();
+ expect(normalizeEmail("user@example.com")).toBe("user@example.com");
+ });
+
+ it("rejects internal whitespace and overlong values", () => {
+ expect(normalizeEmail("us er@example.com")).toBeNull();
+ expect(normalizeEmail("a".repeat(250) + "@example.com")).toBeNull();
+ });
+});
+
+describe("assertReadOnlySelect", () => {
+ it("accepts a plain SELECT", () => {
+ expect(assertReadOnlySelect("select email from users")).toBeNull();
+ });
+
+ it("accepts a trailing semicolon and a CTE", () => {
+ expect(assertReadOnlySelect("select email from users;")).toBeNull();
+ expect(
+ assertReadOnlySelect("with a as (select email from users) select email from a"),
+ ).toBeNull();
+ });
+
+ it("rejects empty queries", () => {
+ expect(assertReadOnlySelect(" ")).toMatch(/empty/i);
+ });
+
+ it("rejects multiple statements", () => {
+ expect(assertReadOnlySelect("select 1; drop table users")).toMatch(/single statement/i);
+ });
+
+ it("rejects non-SELECT statements", () => {
+ expect(assertReadOnlySelect("update users set email='x'")).toMatch(/SELECT/i);
+ });
+
+ it("rejects write/DDL keywords even inside a SELECT", () => {
+ expect(assertReadOnlySelect("select email from users where x in (delete from t)")).toMatch(
+ /only read/i,
+ );
+ expect(assertReadOnlySelect("select email from users; pragma table_info(users)")).not.toBeNull();
+ });
+});