import { CLOSE_ADDRESS_SCHEMA, CLOSE_CONTACT_SCHEMA, CLOSE_ACTIVITY_SCHEMA, CLOSE_LEAD_SCHEMA, CLOSE_OPPORTUNITY_SCHEMA, CLOSE_OPPORTUNITY_PAGE_SCHEMA, CLOSE_VALUE_SCHEMA, closePageSchema, } from "@automate.ax/integration-contracts/close" import * as z from "zod" import { defineAction } from "../../../automation/actions" import { branch } from "../../../automation/signal-operators" import type { DurationString, Signal, } from "../../../automation/signal-protocol" import { delay } from "../../core/delay" import { CLOSE_ACCOUNT, getCloseApi, CLOSE_DATE_INPUT_SCHEMA, CLOSE_ID_SCHEMA, CLOSE_PAGINATION_SCHEMA, hasCloseUpdate, } from "./lib" const CUSTOM_FIELDS_SCHEMA = z.record(z.string().startsWith("cf_"), z.json()) const CLOSE_LEAD_MERGE_MAX_POLLS_SCHEMA = z .number() .int() .min(1) .max(120) .default(60) const CLOSE_LEAD_MERGE_INPUT_SCHEMA = z .object({ destinationLeadId: CLOSE_ID_SCHEMA, sourceLeadId: CLOSE_ID_SCHEMA, }) .refine( ({ destinationLeadId, sourceLeadId }) => destinationLeadId !== sourceLeadId, { message: "Source and destination leads must be different.", path: ["sourceLeadId"], }, ) const CLOSE_LEAD_MERGE_REFERENCE_SCHEMA = z.object({ destinationLeadId: z.string(), sourceLeadId: z.string(), startedAt: z.iso.datetime(), }) const CLOSE_LEAD_MERGE_STATUS_SCHEMA = CLOSE_LEAD_MERGE_REFERENCE_SCHEMA.extend( { activityId: z.string().nullable(), errorMessage: z.string().nullable(), mergeStatus: z.enum(["pending", "in_progress", "completed", "errored"]), }, ) const CLOSE_LEAD_MERGE_OUTPUT_SCHEMA = z.object({ destinationLeadId: z.string(), merged: z.literal(true), sourceLeadId: z.string(), }) /** LeadMerge member extracted from the provider's broad activity union. */ type CloseLeadMergeActivity = Extract< z.output, { type: "LeadMerge" } > const CONTACT_EMAIL_INPUT_SCHEMA = z.object({ email: z.email(), type: z.string().trim().min(1).optional(), }) const CONTACT_PHONE_INPUT_SCHEMA = z.object({ phone: z.string().trim().min(1).nullable().optional(), type: z.string().trim().min(1).optional(), }) const CONTACT_URL_INPUT_SCHEMA = z.object({ type: z.string().trim().min(1).optional(), url: z.url(), }) const NESTED_CONTACT_INPUT_SCHEMA = z.object({ emails: CONTACT_EMAIL_INPUT_SCHEMA.array().optional(), name: z.string().nullable().optional(), phones: CONTACT_PHONE_INPUT_SCHEMA.array().optional(), title: z.string().nullable().optional(), urls: CONTACT_URL_INPUT_SCHEMA.array().optional(), }) const ADDRESS_INPUT_SCHEMA = CLOSE_ADDRESS_SCHEMA.partial() const LEAD_FIELDS = { addresses: ADDRESS_INPUT_SCHEMA.array().optional(), contacts: NESTED_CONTACT_INPUT_SCHEMA.array().optional(), customFields: CUSTOM_FIELDS_SCHEMA.optional(), description: z.string().nullable().optional(), name: z.string().nullable().optional(), source: z.string().nullable().optional(), status: z.string().trim().min(1).optional(), statusId: CLOSE_ID_SCHEMA.optional(), url: z.string().nullable().optional(), } const CONTACT_FIELDS = { customFields: CUSTOM_FIELDS_SCHEMA.optional(), emails: CONTACT_EMAIL_INPUT_SCHEMA.array().max(100).nullable().optional(), leadId: CLOSE_ID_SCHEMA.nullable().optional(), name: z.string().max(1_000).nullable().optional(), phones: CONTACT_PHONE_INPUT_SCHEMA.array().max(100).nullable().optional(), timezone: z.string().nullable().optional(), title: z.string().max(1_000).nullable().optional(), urls: CONTACT_URL_INPUT_SCHEMA.array().max(100).nullable().optional(), } const OPPORTUNITY_ATTACHMENT_SCHEMA = z.object({ contentType: z.string().optional(), filename: z.string().trim().min(1), url: z.url().startsWith("https://app.close.com/go/file/"), }) const CREATE_OPPORTUNITY_FIELDS = { attachments: OPPORTUNITY_ATTACHMENT_SCHEMA.array().nullable().optional(), confidence: z.number().int().min(0).max(100).nullable().optional(), contactId: CLOSE_ID_SCHEMA.nullable().optional(), createdBy: CLOSE_ID_SCHEMA.nullable().optional(), customFields: CUSTOM_FIELDS_SCHEMA.optional(), dateCreated: CLOSE_DATE_INPUT_SCHEMA.nullable().optional(), dateWon: CLOSE_DATE_INPUT_SCHEMA.nullable().optional(), leadId: CLOSE_ID_SCHEMA.nullable().optional(), note: z.string().nullable().optional(), noteHtml: z.string().nullable().optional(), pipelineId: CLOSE_ID_SCHEMA.nullable().optional(), statusId: CLOSE_ID_SCHEMA.nullable().optional(), userId: CLOSE_ID_SCHEMA.nullable().optional(), value: z.number().int().nullable().optional(), valuePeriod: z.enum(["annual", "monthly", "one_time"]).nullable().optional(), } const UPDATE_OPPORTUNITY_FIELDS = { attachments: OPPORTUNITY_ATTACHMENT_SCHEMA.array().optional(), confidence: z.number().int().min(0).max(100).optional(), contactId: CLOSE_ID_SCHEMA.nullable().optional(), customFields: CUSTOM_FIELDS_SCHEMA.optional(), dateWon: CLOSE_DATE_INPUT_SCHEMA.nullable().optional(), note: z.string().optional(), noteHtml: z.string().nullable().optional(), status: z.string().trim().min(1).optional(), statusId: CLOSE_ID_SCHEMA.optional(), userId: CLOSE_ID_SCHEMA.nullable().optional(), value: z.number().int().optional(), valuePeriod: z.enum(["annual", "monthly", "one_time"]).optional(), } /** Lists Close leads with offset pagination. */ export const listCloseLeads = defineAction("List Close leads") .describe("Lists Close leads with offset pagination.") .account(CLOSE_ACCOUNT) .input(z.object(CLOSE_PAGINATION_SCHEMA)) .output(closePageSchema(CLOSE_LEAD_SCHEMA)) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("lead/", { query: { _limit: input.limit, _skip: input.skip }, responseSchema: closePageSchema(CLOSE_LEAD_SCHEMA), }), ) /** Retrieves one Close lead. */ export const getCloseLead = defineAction("Get Close lead") .describe("Retrieves one Close lead by ID.") .account(CLOSE_ACCOUNT) .input(z.object({ leadId: CLOSE_ID_SCHEMA })) .output(CLOSE_LEAD_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request( `lead/${encodeURIComponent(input.leadId)}/`, { responseSchema: CLOSE_LEAD_SCHEMA }, ), ) /** Creates a Close lead with nested contacts and custom fields. */ export const createCloseLead = defineAction("Create Close lead") .describe("Creates a Close lead with addresses, contacts, and custom fields.") .account(CLOSE_ACCOUNT) .input( z .object(LEAD_FIELDS) .refine(({ status, statusId }) => !(status && statusId), { message: "Provide status or statusId, not both.", path: ["statusId"], }), ) .output(CLOSE_LEAD_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("lead/", { body: input, method: "POST", responseSchema: CLOSE_LEAD_SCHEMA, }), ) /** Updates a Close lead. */ export const updateCloseLead = defineAction("Update Close lead") .describe("Updates mutable fields on one Close lead.") .account(CLOSE_ACCOUNT) .input( z .object({ leadId: CLOSE_ID_SCHEMA, ...LEAD_FIELDS }) .refine( (input) => hasCloseUpdate(input, "leadId"), "Provide a lead field to update.", ) .refine(({ status, statusId }) => !(status && statusId), { message: "Provide status or statusId, not both.", path: ["statusId"], }), ) .output(CLOSE_LEAD_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input: { leadId, ...body } }) => getCloseApi(account.secret).request(`lead/${encodeURIComponent(leadId)}/`, { body, method: "PUT", responseSchema: CLOSE_LEAD_SCHEMA, }), ) const startCloseLeadMerge = defineAction("Start Close lead merge") .describe("Starts merging a source lead into a destination lead.") .account(CLOSE_ACCOUNT) .input(CLOSE_LEAD_MERGE_INPUT_SCHEMA) .output(CLOSE_LEAD_MERGE_REFERENCE_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { // Capture before the mutation so polling cannot match an older merge. const startedAt = new Date().toISOString() await getCloseApi(account.secret).request("lead/merge/", { body: { destination: input.destinationLeadId, source: input.sourceLeadId, }, method: "POST", responseSchema: CLOSE_VALUE_SCHEMA, }) return { destinationLeadId: input.destinationLeadId, sourceLeadId: input.sourceLeadId, startedAt, } }) const getCloseLeadMergeStatus = defineAction("Get Close lead merge status") .describe("Gets the provider status for a started Close lead merge.") .account(CLOSE_ACCOUNT) .input(CLOSE_LEAD_MERGE_REFERENCE_SCHEMA) .output(CLOSE_LEAD_MERGE_STATUS_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const activity = ( await getCloseApi(account.secret).request("activity/lead_merge/", { query: { _limit: 100, date_created__gte: input.startedAt, }, responseSchema: closePageSchema(CLOSE_ACTIVITY_SCHEMA), }) ).data .filter( (candidate): candidate is CloseLeadMergeActivity => candidate.type === "LeadMerge" && candidate.destinationLeadId === input.destinationLeadId && candidate.sourceLeadId === input.sourceLeadId, ) .toSorted( (left, right) => right.createdAt.getTime() - left.createdAt.getTime(), )[0] return { ...input, activityId: activity?.id ?? null, errorMessage: activity?.errorMessage ?? null, mergeStatus: activity?.mergeStatus ?? ("pending" as const), } }) /** Input for starting and durably waiting for a Close lead merge. */ export type CloseLeadMergeInput = Extract< Parameters[0], { destinationLeadId: unknown } > & { /** Maximum provider status checks. Defaults to 60. */ maxPolls?: number /** Durable delay between checks. Defaults to `"2s"`. */ pollInterval?: DurationString | number } /** Account selection for a composed Close lead merge. */ export type CloseLeadMergeAccountOptions = NonNullable< Parameters[1] > /** * Merges a source Close lead into a destination and waits for completion. * * Close accepts this mutation asynchronously. Each status check runs as a * separate action after a durable delay, so retrying a check cannot replay the * merge request. * * @param input - Lead IDs and optional bounded polling configuration. * @param options - Optional Close account selection. */ export function mergeCloseLeads( input: CloseLeadMergeInput, options?: CloseLeadMergeAccountOptions, ) { const { maxPolls, pollInterval, ...merge } = input const started = startCloseLeadMerge(merge, options) return pollCloseLeadMerge( CLOSE_LEAD_MERGE_MAX_POLLS_SCHEMA.parse(maxPolls), options, getCloseLeadMergeStatus( { destinationLeadId: started.destinationLeadId, sourceLeadId: started.sourceLeadId, startedAt: started.startedAt, }, options, ), pollInterval ?? "2s", 1, ) } /** * Plans one Close merge status branch and its next durable check. * * @param maxPolls - Total provider checks allowed. * @param options - Close account selection used for every check. * @param status - Current provider merge status. * @param pollInterval - Durable delay between status checks. * @param poll - One-based index of the current status check. */ function pollCloseLeadMerge( maxPolls: number, options: CloseLeadMergeAccountOptions | undefined, status: ReturnType, pollInterval: DurationString | number, poll: number, ): Signal> { return branch( status.mergeStatus.transform( (mergeStatus) => mergeStatus === "pending" || mergeStatus === "in_progress", ), () => { if (poll === maxPolls) { return status.transform(({ destinationLeadId, sourceLeadId }) => { throw new Error( `Close lead merge from ${sourceLeadId} into ${destinationLeadId} remained pending after ${maxPolls} checks at ${pollInterval} intervals.`, ) }) } return delay(pollInterval, () => pollCloseLeadMerge( maxPolls, options, getCloseLeadMergeStatus( { destinationLeadId: status.destinationLeadId, sourceLeadId: status.sourceLeadId, startedAt: status.startedAt, }, options, ), pollInterval, poll + 1, ), ) }, () => status.transform( ({ destinationLeadId, errorMessage, mergeStatus, sourceLeadId }) => { if (mergeStatus === "errored") { throw new Error( `Close lead merge from ${sourceLeadId} into ${destinationLeadId} failed: ${errorMessage ?? "Close did not provide an error message."}`, ) } return { destinationLeadId, merged: true as const, sourceLeadId, } }, ), ) } /** Permanently deletes a Close lead. */ export const deleteCloseLead = defineAction("Delete Close lead") .describe("Permanently deletes a Close lead and its child records.") .account(CLOSE_ACCOUNT) .input(z.object({ leadId: CLOSE_ID_SCHEMA })) .output(z.object({ deleted: z.literal(true), leadId: z.string() })) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await getCloseApi(account.secret).request( `lead/${encodeURIComponent(input.leadId)}/`, { method: "DELETE", responseSchema: CLOSE_VALUE_SCHEMA }, ) return { deleted: true as const, leadId: input.leadId } }) /** Lists Close contacts with optional lead filtering. */ export const listCloseContacts = defineAction("List Close contacts") .describe("Lists Close contacts, optionally within one lead.") .account(CLOSE_ACCOUNT) .input( z.object({ ...CLOSE_PAGINATION_SCHEMA, leadId: CLOSE_ID_SCHEMA.optional(), }), ) .output(closePageSchema(CLOSE_CONTACT_SCHEMA)) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("contact/", { query: { _limit: input.limit, _skip: input.skip, lead_id: input.leadId, }, responseSchema: closePageSchema(CLOSE_CONTACT_SCHEMA), }), ) /** Retrieves one Close contact. */ export const getCloseContact = defineAction("Get Close contact") .describe("Retrieves one Close contact by ID.") .account(CLOSE_ACCOUNT) .input(z.object({ contactId: CLOSE_ID_SCHEMA })) .output(CLOSE_CONTACT_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request( `contact/${encodeURIComponent(input.contactId)}/`, { responseSchema: CLOSE_CONTACT_SCHEMA }, ), ) /** Creates a Close contact. */ export const createCloseContact = defineAction("Create Close contact") .describe("Creates a contact on a Close lead.") .account(CLOSE_ACCOUNT) .input(z.object(CONTACT_FIELDS)) .output(CLOSE_CONTACT_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("contact/", { body: input, method: "POST", responseSchema: CLOSE_CONTACT_SCHEMA, }), ) /** Updates a Close contact. */ export const updateCloseContact = defineAction("Update Close contact") .describe("Updates mutable fields on one Close contact.") .account(CLOSE_ACCOUNT) .input( z .object({ contactId: CLOSE_ID_SCHEMA, ...CONTACT_FIELDS }) .refine( (input) => hasCloseUpdate(input, "contactId"), "Provide a contact field to update.", ), ) .output(CLOSE_CONTACT_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input: { contactId, ...body } }) => getCloseApi(account.secret).request( `contact/${encodeURIComponent(contactId)}/`, { body, method: "PUT", responseSchema: CLOSE_CONTACT_SCHEMA }, ), ) /** Permanently deletes a Close contact. */ export const deleteCloseContact = defineAction("Delete Close contact") .describe("Permanently deletes one Close contact.") .account(CLOSE_ACCOUNT) .input(z.object({ contactId: CLOSE_ID_SCHEMA })) .output(z.object({ contactId: z.string(), deleted: z.literal(true) })) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await getCloseApi(account.secret).request( `contact/${encodeURIComponent(input.contactId)}/`, { method: "DELETE", responseSchema: CLOSE_VALUE_SCHEMA }, ) return { contactId: input.contactId, deleted: true as const } }) /** Lists Close opportunities with common CRM filters. */ export const listCloseOpportunities = defineAction("List Close opportunities") .describe("Lists Close opportunities with lead, status, and owner filters.") .account(CLOSE_ACCOUNT) .input( z.object({ ...CLOSE_PAGINATION_SCHEMA, leadId: CLOSE_ID_SCHEMA.optional(), statusId: CLOSE_ID_SCHEMA.optional(), statusType: z.enum(["active", "lost", "won"]).optional(), userId: CLOSE_ID_SCHEMA.optional(), }), ) .output(CLOSE_OPPORTUNITY_PAGE_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("opportunity/", { query: { _limit: input.limit, _skip: input.skip, lead_id: input.leadId, status_id: input.statusId, status_type: input.statusType, user_id: input.userId, }, responseSchema: CLOSE_OPPORTUNITY_PAGE_SCHEMA, }), ) /** Retrieves one Close opportunity. */ export const getCloseOpportunity = defineAction("Get Close opportunity") .describe("Retrieves one Close opportunity by ID.") .account(CLOSE_ACCOUNT) .input(z.object({ opportunityId: CLOSE_ID_SCHEMA })) .output(CLOSE_OPPORTUNITY_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request( `opportunity/${encodeURIComponent(input.opportunityId)}/`, { responseSchema: CLOSE_OPPORTUNITY_SCHEMA }, ), ) /** Creates a Close opportunity. */ export const createCloseOpportunity = defineAction("Create Close opportunity") .describe("Creates an opportunity on a Close lead and pipeline status.") .account(CLOSE_ACCOUNT) .input(z.object(CREATE_OPPORTUNITY_FIELDS)) .output(CLOSE_OPPORTUNITY_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => getCloseApi(account.secret).request("opportunity/", { body: input, method: "POST", responseSchema: CLOSE_OPPORTUNITY_SCHEMA, }), ) /** Updates a Close opportunity. */ export const updateCloseOpportunity = defineAction("Update Close opportunity") .describe("Updates deal value, ownership, notes, or pipeline status.") .account(CLOSE_ACCOUNT) .input( z .object({ opportunityId: CLOSE_ID_SCHEMA, ...UPDATE_OPPORTUNITY_FIELDS, }) .refine( (input) => hasCloseUpdate(input, "opportunityId"), "Provide an opportunity field to update.", ) .refine(({ status, statusId }) => !(status && statusId), { message: "Provide status or statusId, not both.", path: ["statusId"], }), ) .output(CLOSE_OPPORTUNITY_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input: { opportunityId, ...body } }) => getCloseApi(account.secret).request( `opportunity/${encodeURIComponent(opportunityId)}/`, { body, method: "PUT", responseSchema: CLOSE_OPPORTUNITY_SCHEMA }, ), ) /** Permanently deletes a Close opportunity. */ export const deleteCloseOpportunity = defineAction("Delete Close opportunity") .describe("Permanently deletes one Close opportunity.") .account(CLOSE_ACCOUNT) .input(z.object({ opportunityId: CLOSE_ID_SCHEMA })) .output(z.object({ deleted: z.literal(true), opportunityId: z.string() })) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await getCloseApi(account.secret).request( `opportunity/${encodeURIComponent(input.opportunityId)}/`, { method: "DELETE", responseSchema: CLOSE_VALUE_SCHEMA }, ) return { deleted: true as const, opportunityId: input.opportunityId } })