Skip to content
Open
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
16 changes: 16 additions & 0 deletions apps/api/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -8,3 +8,19 @@ RESEND_API_KEY=re_xxxxx
RESEND_FROM_EMAIL=Verio<onboarding@xxxxx>
GOOGLE_CLIENT_ID=xxxxx
GOOGLE_CLIENT_SECRET=xxxxx

# Import pipeline
# Leave false to run every connector against built-in fixtures. The full pipeline —
# extraction, mapping, matching, loading — works end to end with no provider credentials.
INTEGRATION_LIVE_FETCH_ENABLED=false
# Encrypts stored OAuth and API tokens. Falls back to BETTER_AUTH_SECRET when unset.
INTEGRATION_TOKEN_SECRET=<generate with: openssl rand -base64 32>

# Outlook / Microsoft Graph (contacts, sent mail, calendar through one grant)
MICROSOFT_CLIENT_ID=xxxxx
MICROSOFT_CLIENT_SECRET=xxxxx
MICROSOFT_TENANT_ID=common

# Use https://eu.posthog.com for EU projects, or a self-hosted host
POSTHOG_API_HOST=https://us.posthog.com
CALENDLY_API_HOST=https://api.calendly.com
1 change: 1 addition & 0 deletions apps/api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
"build": "tsc",
"start": "bun run src/server.ts",
"worker:sequences": "bun run src/workers/sequence-worker.ts",
"worker:imports": "bun run src/workers/import-worker.ts",
"check-types": "tsc --noEmit && tsc --noEmit -p test/tsconfig.json",
"lint": "eslint .",
"test": "TZ=UTC vitest run",
Expand Down
7 changes: 7 additions & 0 deletions apps/api/src/config/env.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,13 @@ const rawEnv = {
GOOGLE_CLIENT_SECRET: process.env.GOOGLE_CLIENT_SECRET,
SEQUENCE_LIVE_SEND_ENABLED: process.env.SEQUENCE_LIVE_SEND_ENABLED,
GROQ_API_KEY: process.env.GROQ_API_KEY,
INTEGRATION_LIVE_FETCH_ENABLED: process.env.INTEGRATION_LIVE_FETCH_ENABLED,
INTEGRATION_TOKEN_SECRET: process.env.INTEGRATION_TOKEN_SECRET,
MICROSOFT_CLIENT_ID: process.env.MICROSOFT_CLIENT_ID,
MICROSOFT_CLIENT_SECRET: process.env.MICROSOFT_CLIENT_SECRET,
MICROSOFT_TENANT_ID: process.env.MICROSOFT_TENANT_ID,
POSTHOG_API_HOST: process.env.POSTHOG_API_HOST,
CALENDLY_API_HOST: process.env.CALENDLY_API_HOST,
};

const parsedEnv = apiEnvSchema.safeParse(rawEnv);
Expand Down
185 changes: 185 additions & 0 deletions apps/api/src/controllers/imports.controller.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
import type { Context } from "hono";
import type {
CreateConnectionInput,
ListImportJobsQuery,
ListImportRecordsQuery,
StartImportJobInput,
UpdateConnectionInput,
UpdateImportJobInput,
WebhookIngestInput,
} from "@workspace/validators/schemas/import";
import type { ImportEntityType } from "@workspace/validators/types/import";
import { STATUS_CODES } from "@/constants/status-codes.js";
import { sendSuccess } from "@/lib/api-response.js";
import { AppError } from "@/lib/app-error.js";
import { getSessionWorkspaceId } from "@/lib/workspace.js";
import {
cancelImportJob,
commitImportJob,
createApiKey,
createConnection,
deleteConnection,
getImportJob,
ingestWebhookRecords,
listApiKeys,
listConnections,
listImportJobs,
listImportRecords,
previewSourceFields,
resolveApiKey,
revokeApiKey,
startImportJob,
updateConnection,
updateImportJob,
} from "@/services/import.service.js";

export async function listConnectionsController(c: Context) {
const workspaceId = getSessionWorkspaceId(c);
const connections = await listConnections(workspaceId);

return sendSuccess(c, { connections }, STATUS_CODES.OK);
}

export async function createConnectionController(c: Context, payload: CreateConnectionInput) {
const workspaceId = getSessionWorkspaceId(c);
const user = c.get("user");
const connection = await createConnection(workspaceId, user.id, payload);

return sendSuccess(c, { connection }, STATUS_CODES.CREATED);
}

export async function updateConnectionController(
c: Context,
id: string,
payload: UpdateConnectionInput,
) {
const workspaceId = getSessionWorkspaceId(c);
const connection = await updateConnection(workspaceId, id, payload);

return sendSuccess(c, { connection }, STATUS_CODES.OK);
}

export async function deleteConnectionController(c: Context, id: string) {
const workspaceId = getSessionWorkspaceId(c);
const connection = await deleteConnection(workspaceId, id);

return sendSuccess(c, { connection }, STATUS_CODES.OK);
}

export async function previewConnectionFieldsController(
c: Context,
id: string,
entityType: ImportEntityType,
) {
const workspaceId = getSessionWorkspaceId(c);
const preview = await previewSourceFields(workspaceId, id, entityType);

return sendSuccess(c, preview, STATUS_CODES.OK);
}

export async function listApiKeysController(c: Context) {
const workspaceId = getSessionWorkspaceId(c);
const apiKeys = await listApiKeys(workspaceId);

return sendSuccess(c, { apiKeys }, STATUS_CODES.OK);
}

export async function createApiKeyController(c: Context, payload: { name: string }) {
const workspaceId = getSessionWorkspaceId(c);
const user = c.get("user");
const result = await createApiKey(workspaceId, user.id, payload.name);

return sendSuccess(
c,
result,
STATUS_CODES.CREATED,
"Store this key now. It cannot be shown again.",
);
}

export async function revokeApiKeyController(c: Context, id: string) {
const workspaceId = getSessionWorkspaceId(c);
const apiKey = await revokeApiKey(workspaceId, id);

return sendSuccess(c, { apiKey }, STATUS_CODES.OK);
}

export async function startImportJobController(c: Context, payload: StartImportJobInput) {
const workspaceId = getSessionWorkspaceId(c);
const user = c.get("user");
const job = await startImportJob(workspaceId, user.id, payload);

return sendSuccess(c, { job }, STATUS_CODES.CREATED);
}

export async function listImportJobsController(c: Context, query: ListImportJobsQuery) {
const workspaceId = getSessionWorkspaceId(c);
const result = await listImportJobs(workspaceId, query);

return sendSuccess(c, result, STATUS_CODES.OK);
}

export async function getImportJobController(c: Context, id: string) {
const workspaceId = getSessionWorkspaceId(c);
const job = await getImportJob(workspaceId, id);

return sendSuccess(c, { job }, STATUS_CODES.OK);
}

export async function listImportRecordsController(
c: Context,
id: string,
query: ListImportRecordsQuery,
) {
const workspaceId = getSessionWorkspaceId(c);
const result = await listImportRecords(workspaceId, id, query);

return sendSuccess(c, result, STATUS_CODES.OK);
}

export async function updateImportJobController(
c: Context,
id: string,
payload: UpdateImportJobInput,
) {
const workspaceId = getSessionWorkspaceId(c);
const job = await updateImportJob(workspaceId, id, payload);

return sendSuccess(c, { job }, STATUS_CODES.OK);
}

export async function commitImportJobController(c: Context, id: string) {
const workspaceId = getSessionWorkspaceId(c);
const job = await commitImportJob(workspaceId, id);

return sendSuccess(c, { job }, STATUS_CODES.ACCEPTED);
}

export async function cancelImportJobController(c: Context, id: string) {
const workspaceId = getSessionWorkspaceId(c);
const job = await cancelImportJob(workspaceId, id);

return sendSuccess(c, { job }, STATUS_CODES.OK);
}

/**
* Unauthenticated by session — this is the endpoint Zapier, Make, n8n, and custom scripts
* push to, so it authenticates with a workspace API key instead.
*/
export async function ingestWebhookController(c: Context, payload: WebhookIngestInput) {
const header = c.req.header("authorization") ?? "";
const token = header.toLowerCase().startsWith("bearer ") ? header.slice(7).trim() : "";

if (token === "") {
throw new AppError("Missing API key", STATUS_CODES.UNAUTHORIZED);
}

const apiKey = await resolveApiKey(token);
if (!apiKey) {
throw new AppError("Invalid or revoked API key", STATUS_CODES.UNAUTHORIZED);
}

const result = await ingestWebhookRecords(apiKey.workspaceId, payload);

return sendSuccess(c, result, STATUS_CODES.ACCEPTED);
}
126 changes: 126 additions & 0 deletions apps/api/src/db/drizzle/0004_public_mesmero.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
CREATE TYPE "public"."import_connection_status" AS ENUM('connected', 'reconnect_required', 'disconnected');--> statement-breakpoint
CREATE TYPE "public"."import_entity_type" AS ENUM('person', 'org');--> statement-breakpoint
CREATE TYPE "public"."import_job_status" AS ENUM('pending', 'extracting', 'ready_for_review', 'loading', 'completed', 'failed', 'canceled');--> statement-breakpoint
CREATE TYPE "public"."import_match_reason" AS ENUM('external_identity', 'email', 'domain', 'name', 'none');--> statement-breakpoint
CREATE TYPE "public"."import_provider" AS ENUM('csv', 'webhook', 'gmail', 'google_calendar', 'calendly', 'google_sheets', 'posthog', 'outlook');--> statement-breakpoint
CREATE TYPE "public"."import_record_status" AS ENUM('pending', 'valid', 'invalid', 'loaded', 'skipped', 'duplicate');--> statement-breakpoint
ALTER TYPE "public"."people_source" ADD VALUE 'import';--> statement-breakpoint
CREATE TABLE "external_identities" (
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
"workspace_id" uuid NOT NULL,
"provider" "import_provider" NOT NULL,
"external_id" varchar(255) NOT NULL,
"entity_type" "import_entity_type" NOT NULL,
"person_id" uuid,
"org_id" uuid,
"profile" jsonb DEFAULT '{}'::jsonb NOT NULL,
"last_seen_at" timestamp DEFAULT now() NOT NULL,
"created_at" timestamp DEFAULT now() NOT NULL,
"updated_at" timestamp DEFAULT now() NOT NULL,
CONSTRAINT "external_identities_workspace_provider_external_unique" UNIQUE("workspace_id","provider","external_id")
);
--> statement-breakpoint
CREATE TABLE "import_jobs" (
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
"workspace_id" uuid NOT NULL,
"connection_id" uuid,
"created_by_id" uuid,
"provider" "import_provider" NOT NULL,
"entity_type" "import_entity_type" DEFAULT 'person' NOT NULL,
"status" "import_job_status" DEFAULT 'pending' NOT NULL,
"mapping" jsonb DEFAULT '{"fields":[]}'::jsonb NOT NULL,
"options" jsonb NOT NULL,
"source_config" jsonb DEFAULT '{}'::jsonb NOT NULL,
"stats" jsonb NOT NULL,
"cursor" jsonb,
"attempts" integer DEFAULT 0 NOT NULL,
"next_attempt_at" timestamp,
"started_at" timestamp,
"completed_at" timestamp,
"last_error" text,
"created_at" timestamp DEFAULT now() NOT NULL,
"updated_at" timestamp DEFAULT now() NOT NULL
);
--> statement-breakpoint
CREATE TABLE "import_records" (
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
"workspace_id" uuid NOT NULL,
"job_id" uuid NOT NULL,
"external_id" varchar(255),
"row_number" integer NOT NULL,
"raw" jsonb NOT NULL,
"normalized" jsonb,
"status" "import_record_status" DEFAULT 'pending' NOT NULL,
"match_reason" "import_match_reason" DEFAULT 'none' NOT NULL,
"match_person_id" uuid,
"match_org_id" uuid,
"errors" jsonb DEFAULT '[]'::jsonb NOT NULL,
"created_at" timestamp DEFAULT now() NOT NULL,
"updated_at" timestamp DEFAULT now() NOT NULL,
CONSTRAINT "import_records_job_row_unique" UNIQUE("job_id","row_number")
);
--> statement-breakpoint
CREATE TABLE "integration_api_keys" (
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
"workspace_id" uuid NOT NULL,
"created_by_id" uuid,
"name" varchar(255) NOT NULL,
"token_hash" varchar(128) NOT NULL,
"token_prefix" varchar(16) NOT NULL,
"last_used_at" timestamp,
"revoked_at" timestamp,
"created_at" timestamp DEFAULT now() NOT NULL,
"updated_at" timestamp DEFAULT now() NOT NULL,
CONSTRAINT "integration_api_keys_token_hash_unique" UNIQUE("token_hash")
);
--> statement-breakpoint
CREATE TABLE "integration_connections" (
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
"workspace_id" uuid NOT NULL,
"user_id" uuid NOT NULL,
"provider" "import_provider" NOT NULL,
"display_name" varchar(255) NOT NULL,
"external_account_id" varchar(255),
"status" "import_connection_status" DEFAULT 'connected' NOT NULL,
"granted_scopes" jsonb DEFAULT '[]'::jsonb NOT NULL,
"access_token_encrypted" text NOT NULL,
"refresh_token_encrypted" text,
"token_expires_at" timestamp,
"config" jsonb DEFAULT '{}'::jsonb NOT NULL,
"cursor" jsonb,
"last_sync_at" timestamp,
"last_error" text,
"created_at" timestamp DEFAULT now() NOT NULL,
"updated_at" timestamp DEFAULT now() NOT NULL,
CONSTRAINT "integration_connections_workspace_provider_account_unique" UNIQUE("workspace_id","provider","external_account_id")
);
--> statement-breakpoint
ALTER TABLE "external_identities" ADD CONSTRAINT "external_identities_workspace_id_workspaces_id_fk" FOREIGN KEY ("workspace_id") REFERENCES "public"."workspaces"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "external_identities" ADD CONSTRAINT "external_identities_person_id_people_id_fk" FOREIGN KEY ("person_id") REFERENCES "public"."people"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "external_identities" ADD CONSTRAINT "external_identities_org_id_org_id_fk" FOREIGN KEY ("org_id") REFERENCES "public"."org"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_jobs" ADD CONSTRAINT "import_jobs_workspace_id_workspaces_id_fk" FOREIGN KEY ("workspace_id") REFERENCES "public"."workspaces"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_jobs" ADD CONSTRAINT "import_jobs_connection_id_integration_connections_id_fk" FOREIGN KEY ("connection_id") REFERENCES "public"."integration_connections"("id") ON DELETE set null ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_jobs" ADD CONSTRAINT "import_jobs_created_by_id_user_id_fk" FOREIGN KEY ("created_by_id") REFERENCES "public"."user"("id") ON DELETE set null ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_records" ADD CONSTRAINT "import_records_workspace_id_workspaces_id_fk" FOREIGN KEY ("workspace_id") REFERENCES "public"."workspaces"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_records" ADD CONSTRAINT "import_records_job_id_import_jobs_id_fk" FOREIGN KEY ("job_id") REFERENCES "public"."import_jobs"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_records" ADD CONSTRAINT "import_records_match_person_id_people_id_fk" FOREIGN KEY ("match_person_id") REFERENCES "public"."people"("id") ON DELETE set null ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "import_records" ADD CONSTRAINT "import_records_match_org_id_org_id_fk" FOREIGN KEY ("match_org_id") REFERENCES "public"."org"("id") ON DELETE set null ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "integration_api_keys" ADD CONSTRAINT "integration_api_keys_workspace_id_workspaces_id_fk" FOREIGN KEY ("workspace_id") REFERENCES "public"."workspaces"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "integration_api_keys" ADD CONSTRAINT "integration_api_keys_created_by_id_user_id_fk" FOREIGN KEY ("created_by_id") REFERENCES "public"."user"("id") ON DELETE set null ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "integration_connections" ADD CONSTRAINT "integration_connections_workspace_id_workspaces_id_fk" FOREIGN KEY ("workspace_id") REFERENCES "public"."workspaces"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "integration_connections" ADD CONSTRAINT "integration_connections_user_id_user_id_fk" FOREIGN KEY ("user_id") REFERENCES "public"."user"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
CREATE INDEX "external_identities_workspace_id_idx" ON "external_identities" USING btree ("workspace_id");--> statement-breakpoint
CREATE INDEX "external_identities_person_id_idx" ON "external_identities" USING btree ("person_id");--> statement-breakpoint
CREATE INDEX "external_identities_org_id_idx" ON "external_identities" USING btree ("org_id");--> statement-breakpoint
CREATE INDEX "import_jobs_workspace_id_idx" ON "import_jobs" USING btree ("workspace_id");--> statement-breakpoint
CREATE INDEX "import_jobs_connection_id_idx" ON "import_jobs" USING btree ("connection_id");--> statement-breakpoint
CREATE INDEX "import_jobs_status_next_attempt_idx" ON "import_jobs" USING btree ("status","next_attempt_at");--> statement-breakpoint
CREATE INDEX "import_jobs_provider_idx" ON "import_jobs" USING btree ("provider");--> statement-breakpoint
CREATE INDEX "import_records_workspace_id_idx" ON "import_records" USING btree ("workspace_id");--> statement-breakpoint
CREATE INDEX "import_records_job_status_idx" ON "import_records" USING btree ("job_id","status");--> statement-breakpoint
CREATE INDEX "import_records_external_id_idx" ON "import_records" USING btree ("external_id");--> statement-breakpoint
CREATE INDEX "integration_api_keys_workspace_id_idx" ON "integration_api_keys" USING btree ("workspace_id");--> statement-breakpoint
CREATE INDEX "integration_connections_workspace_id_idx" ON "integration_connections" USING btree ("workspace_id");--> statement-breakpoint
CREATE INDEX "integration_connections_user_id_idx" ON "integration_connections" USING btree ("user_id");--> statement-breakpoint
CREATE INDEX "integration_connections_provider_idx" ON "integration_connections" USING btree ("provider");--> statement-breakpoint
CREATE INDEX "integration_connections_status_idx" ON "integration_connections" USING btree ("status");
Loading