import errors from '@tryghost/errors'; import tpl from '@tryghost/tpl'; import crypto from 'node:crypto'; import ObjectId from 'bson-objectid'; import {dequal} from 'dequal'; import {type Knex} from 'knex'; // @ts-expect-error This module currently lacks type definitions. import lexicalLib from '../../lib/lexical'; import urlUtils from '../../../shared/url-utils'; import {DEFAULT_EMAIL_DESIGN_SETTING_SLUG, MEMBER_WELCOME_EMAIL_SLUGS} from '../member-welcome-emails/constants'; import type { AutomatedEmailEvents, Automation, AutomationAction, AutomationEdge, AutomationEmailStats, AutomationSummary, AutomationStepTerminalStatus, AutomationStepToRun, AutomationsRepository, EditAutomationData, Page } from './automations-repository'; import {fromDatabaseDate, toDatabaseDate, type DatabaseDate} from './database-date'; import {getStaleLockCutoff} from './stale-lock-cutoff'; import type {ExclusifyUnion, ReadonlyDeep} from 'type-fest'; const HOUR_MS = 60 * 60 * 1000; const DEFAULT_WELCOME_EMAIL_AUTOMATIONS = [{ name: 'Free member welcome flow', slug: MEMBER_WELCOME_EMAIL_SLUGS.free }, { name: 'Paid member welcome flow', slug: MEMBER_WELCOME_EMAIL_SLUGS.paid }]; const messages = { invalidAutomationActionRevision: 'Automation action "{actionId}" of type "{actionType}" is missing required revision field "{field}".', conflictingAutomationActionId: 'Automation action "{actionId}" already exists and cannot be inserted.', conflictingAutomationActionType: 'Automation action "{actionId}" already exists with a different type.', defaultEmailDesignSettingNotFound: 'Default automated email design setting not found.' }; const DEFAULT_EMAIL_DESIGN_SETTING_REFERENCE = DEFAULT_EMAIL_DESIGN_SETTING_SLUG; interface AutomationRow { id: string; slug: string; name: string; status: string; created_at: DatabaseDate; updated_at: DatabaseDate; } interface ActionRow { id: string; type: 'wait' | 'send_email'; wait_hours: number | null; email_subject: string | null; email_lexical: string | null; email_design_setting_id: string | null; } type ActionStatsRow = { action_id: string; email_clicked_count: number | null; email_sent_count: number | null; email_opened_count: number | null; }; type ActionLinkRow = { url: string; clicked_count: string | number; }; type ActionRevisionRow = { action_id: string; created_at: DatabaseDate; wait_hours: number | null; email_subject: string | null; email_lexical: string | null; email_design_setting_id: string | null; }; interface EdgeRow { source_action_id: string; target_action_id: string; } type NextActionRevisionRow = { automation_id: string; action_id: string; automation_action_revision_id: string; type: 'wait' | 'send_email'; wait_hours: number | null; }; type StepToRunRow = { id: string; locked_by: string; automation_run_id: string; automation_id: string; automation_slug: string; automation_status: 'inactive' | 'active'; member_id: string | null; member_email: string; action_id: string; automation_action_revision_id: string; type: string; ready_at: DatabaseDate; step_attempts: number; wait_hours: number | null; email_subject: string | null; email_lexical: string | null; email_design_setting_id: string | null; }; type WaitActionData = Extract['data']; type SendEmailActionData = Extract['data']; type RevisionDataFor = { [FieldT in keyof ActionDataT]: ActionRevisionRow[FieldT & keyof ActionRevisionRow]; }; type WaitRevisionData = RevisionDataFor; type SendEmailRevisionData = RevisionDataFor; type ActionToInsert = { id: string; created_at: string; updated_at: string; automation_id: string; type: AutomationAction['type']; }; type ActionRevisionToInsert = { actionId: string; action: AutomationAction; createdAt: string; }; export function createDatabaseAutomationsRepository({ knex, fakeWaitHoursMultiplier }: { knex: Knex; fakeWaitHoursMultiplier: number | null; }): AutomationsRepository { return { async browse(): Promise> { return await knex.transaction(async (trx) => { await ensureDefaultAutomations(trx); const rows = await loadAutomations(trx); return { data: rows.map(row => buildAutomationSummary(row)), meta: { pagination: buildPagination(rows.length) } }; }); }, async getById(id: string): Promise { return await knex.transaction(async (trx) => { const automation = await loadAutomation(trx, id); if (!automation) { return null; } return await buildAutomation(trx, automation); }); }, async getAutomationActionLinks(automationId, actionId) { const action = await knex('automation_actions') .select('id') .where({ id: actionId, automation_id: automationId }) .whereNull('deleted_at') .first(); if (!action) { return null; } const rows = await knex('redirects as redirects') .countDistinct({clicked_count: 'members_click_events.member_id'}) .select(knex.raw('MIN(??) as ??', ['redirects.to', 'url'])) .innerJoin('automation_action_revisions as revisions', 'revisions.id', 'redirects.automation_action_revision_id') .leftJoin('members_click_events', 'members_click_events.redirect_id', 'redirects.id') .where('revisions.action_id', actionId) .whereNotNull('redirects.to_hash') .groupBy('redirects.to_hash') .orderBy('clicked_count', 'desc') .orderBy('url', 'asc'); return rows.map(row => ({ url: urlUtils.transformReadyToAbsolute(row.url), clicked_count: Number(row.clicked_count) })); }, async edit(id: string, data: EditAutomationData): Promise { return await knex.transaction(async (trx) => { const automation = await loadAutomation(trx, id); if (!automation) { return null; } const now = new Date(); const updatedAutomation = await updateAutomation(trx, { ...automation, status: data.status, updated_at: toDatabaseDate(now) }); await replaceAutomationGraph(trx, updatedAutomation.id, data.actions, data.edges); if (updatedAutomation.status === 'inactive') { await cancelCancelablePendingStepsForAutomation(trx, updatedAutomation.id, now); } return await buildAutomation(trx, updatedAutomation); }); }, async trigger(options: { memberEmail: string; memberId: string; memberStatus: 'free' | 'paid'; }): Promise { return await knex.transaction(trx => trigger(trx, { ...options, fakeWaitHoursMultiplier })); }, async fetchAndLockSteps(limit: number): Promise<{ steps: AutomationStepToRun[], nextStepReadyAt: null | Date; }> { return await knex.transaction(trx => fetchAndLockSteps(trx, limit)); }, async finishStepAndEnqueueNext(step: AutomationStepToRun): Promise { return await knex.transaction(trx => finishStepAndEnqueueNext(trx, { step, fakeWaitHoursMultiplier })); }, async markStepTerminal(step: AutomationStepToRun, status: AutomationStepTerminalStatus): Promise { return await knex.transaction(trx => markStepTerminal(trx, step, status)); }, async retryStep(step: AutomationStepToRun, retryAt: Date): Promise { return await knex.transaction(trx => retryStep(trx, step, retryAt)); }, async recordEmailSent(options): Promise { await knex.transaction(async (trx) => { await trx('automation_action_revisions') .where('id', options.automationActionRevisionId) .update({ email_sent_count: trx.raw('COALESCE(??, 0) + ?', ['email_sent_count', 1]) }); const now = toDatabaseDate(new Date()); await trx('automated_email_recipients').insert({ id: ObjectId().toHexString(), member_id: options.memberId, member_uuid: options.memberUuid, member_email: options.memberEmail, member_name: options.memberName, automation_action_revision_id: options.automationActionRevisionId, automation_run_step_id: options.automationRunStepId, ...(options.mailgunMessageId ? {mailgun_message_id: options.mailgunMessageId} : {}), track_clicks: options.trackClicks, track_opens: options.trackOpens, created_at: now, updated_at: now }); }); }, async getAutomatedEmailRecipientsByMailgunIds(mailgunMessageIds) { if (mailgunMessageIds.length === 0) { return []; } return await knex('automated_email_recipients') .select('id', 'mailgun_message_id', 'automation_action_revision_id') .whereNotNull('automation_action_revision_id') .whereIn('mailgun_message_id', mailgunMessageIds); }, async trackEmailDeliveredAndOpened(eventsByAutomatedEmailRecipientId) { if (eventsByAutomatedEmailRecipientId.size === 0) { return; } await knex.transaction(async (trx) => { const revisionIds = new Set(); for (const {openedAt, automationActionRevisionId} of eventsByAutomatedEmailRecipientId.values()) { if (openedAt) { revisionIds.add(automationActionRevisionId); } } const orderedRevisionIds = await lockActionRevisions(trx, revisionIds); const notYetOpened = await lockNotYetOpened(trx, eventsByAutomatedEmailRecipientId); const newOpensPerRevision = new Map(); for (const [id, {deliveredAt, openedAt, automationActionRevisionId}] of eventsByAutomatedEmailRecipientId) { const updates: Record = {}; if (deliveredAt) { updates.delivered_at = trx.raw('CASE WHEN delivered_at IS NULL OR delivered_at > ? THEN ? ELSE delivered_at END', [deliveredAt, deliveredAt]); } if (openedAt) { updates.opened_at = trx.raw('CASE WHEN opened_at IS NULL OR opened_at > ? THEN ? ELSE opened_at END', [openedAt, openedAt]); } if (Object.keys(updates).length === 0) { continue; } await trx('automated_email_recipients') .where({id}) .update(updates); if (openedAt && notYetOpened.has(id)) { newOpensPerRevision.set( automationActionRevisionId, (newOpensPerRevision.get(automationActionRevisionId) ?? 0) + 1 ); } } for (const id of orderedRevisionIds) { const opens = newOpensPerRevision.get(id); if (!opens) { continue; } await trx('automation_action_revisions') .where({id}) .update({ email_opened_count: trx.raw('COALESCE(email_opened_count, 0) + ?', [opens]) }); } }); }, async trackEmailClicked({automationActionRevisionId, automationRunStepId, memberId, clickedAt}, {transacting} = {}) { const trackClick = async (trx: Knex.Transaction) => { await lockActionRevisions(trx, [automationActionRevisionId]); const recipient = await trx('automated_email_recipients') .select('id', 'clicked_at') .where({ automation_action_revision_id: automationActionRevisionId, automation_run_step_id: automationRunStepId, member_id: memberId, track_clicks: true }) .forUpdate() .first(); if (!recipient || recipient.clicked_at !== null) { return; } const updated = await trx('automated_email_recipients') .where({id: recipient.id}) .whereNull('clicked_at') .update({clicked_at: clickedAt}); if (updated === 0) { return; } await trx('automation_action_revisions') .where({id: automationActionRevisionId}) .update({ email_clicked_count: trx.raw('COALESCE(email_clicked_count, 0) + 1') }); }; if (transacting) { await trackClick(transacting); return; } await knex.transaction(trackClick); } }; } /** * Lock revisions before recipients because inserting a recipient takes a shared * foreign-key lock on its revision. Updating that revision later can deadlock * with another transaction that has already locked the revision and is waiting * for the recipient. Keep multi-revision lock acquisition deterministic. */ async function lockActionRevisions( trx: Knex.Transaction, revisionIds: Iterable ): Promise { const sortedRevisionIds = [...new Set(revisionIds)] .sort((left, right) => left.localeCompare(right)); if (sortedRevisionIds.length > 0) { await trx('automation_action_revisions') .select('id') .whereIn('id', sortedRevisionIds) .forUpdate(); } return sortedRevisionIds; } /** * Which of these recipients have yet to open, and so should count towards their * revision's open count. Locks them for the transaction, so a worker racing on * the same open reads them as opened and doesn't count them a second time. */ async function lockNotYetOpened( trx: Knex.Transaction, eventsByAutomatedEmailRecipientId: ReadonlyDeep> ): Promise> { const ids: string[] = []; for (const [id, {openedAt}] of eventsByAutomatedEmailRecipientId) { if (openedAt) { ids.push(id); } } if (ids.length === 0) { return new Set(); } const rows = await trx('automated_email_recipients') .select('id') .whereIn('id', ids) .whereNull('opened_at') .forUpdate(); return new Set(rows.map(row => row.id)); } async function ensureDefaultAutomations(trx: Knex.Transaction): Promise { for (const defaults of DEFAULT_WELCOME_EMAIL_AUTOMATIONS) { const automation = await ensureAutomation(trx, defaults); await ensureWelcomeEmailAction(trx, automation.id); } } async function ensureAutomation( trx: Knex.Transaction, defaults: Readonly<{name: string; slug: string}> ): Promise { const now = toDatabaseDate(new Date()); const id = ObjectId().toHexString(); await trx('automations') .insert({ id, status: 'inactive', name: defaults.name, slug: defaults.slug, created_at: now, updated_at: now }) .onConflict('slug') .ignore(); return requireAutomation(await loadAutomationBySlug(trx, defaults.slug), defaults.slug); } async function ensureWelcomeEmailAction(trx: Knex.Transaction, automationId: string): Promise { const hasActions = await trx('automation_actions') .where('automation_id', automationId) .whereNull('deleted_at') .first('id'); if (hasActions) { return; } await trx('automations') .select('id') .where('id', automationId) .forUpdate() .first(); const email = await trx('welcome_email_automated_emails') .select('subject', 'lexical', 'email_design_setting_id') .where('welcome_email_automation_id', automationId) .orderBy([ 'created_at', 'id' ]) .first(); if (!email) { return; } const now = toDatabaseDate(new Date()); const actionId = ObjectId().toHexString(); await insertActions(trx, [{ id: actionId, created_at: now, updated_at: now, automation_id: automationId, type: 'send_email' }]); await insertActionRevisions(trx, [{ actionId, action: { id: actionId, type: 'send_email', data: { email_subject: email.subject, email_lexical: email.lexical ?? '', email_design_setting_id: email.email_design_setting_id } }, createdAt: getNextRevisionCreatedAt(null, now) }]); } async function trigger(trx: Knex.Transaction, options: Readonly<{ memberEmail: string; memberId: string; memberStatus: 'free' | 'paid'; fakeWaitHoursMultiplier: number | null; }>): Promise { const { memberEmail, memberId, memberStatus, fakeWaitHoursMultiplier } = options; const firstAction = await findFirstActionRevision(trx, memberStatus); if (!firstAction) { return; } const now = new Date(); const nowString = toDatabaseDate(now); const readyAt = getReadyAtForAction(firstAction, now, fakeWaitHoursMultiplier); const run = { id: ObjectId().toHexString(), created_at: nowString, updated_at: nowString, automation_id: firstAction.automation_id, member_id: memberId, member_email: memberEmail }; await trx('automation_runs').insert(run); await insertRunStep(trx, { automationRunId: run.id, automationActionRevisionId: firstAction.automation_action_revision_id, now, readyAt }); } async function insertRunStep(trx: Knex.Transaction, { automationRunId, automationActionRevisionId, now, readyAt }: ReadonlyDeep<{ automationRunId: string; automationActionRevisionId: string; now: Date; readyAt: Date; }>): Promise { const nowString = toDatabaseDate(now); await trx('automation_run_steps').insert({ id: ObjectId().toHexString(), created_at: nowString, updated_at: nowString, automation_run_id: automationRunId, automation_action_revision_id: automationActionRevisionId, ready_at: toDatabaseDate(readyAt) }); } async function fetchAndLockSteps(trx: Knex.Transaction, limit: number): Promise<{ steps: AutomationStepToRun[], nextStepReadyAt: null | Date; }> { // Two things make this tricky: // // - We want to do row-level locking, so multiple calls don't step on each other. // - We can't `UPDATE` a fixed number of rows. // // To get around these problems, here's what we do: // // 1. Select up to `limit` candidate rows. // 2. Try to lock those rows. // 3. Select any rows we successfully locked. const now = new Date(); const nowString = toDatabaseDate(now); const staleLockCutoff = getStaleLockCutoff(now); const staleLockCutoffString = toDatabaseDate(staleLockCutoff); const lockId = crypto.randomUUID(); // 1. Select up to `limit` candidate rows. const candidates: ReadonlyArray<{id: string}> = await trx('automation_run_steps') .select('id') .where('status', 'pending') .where('ready_at', '<=', nowString) .where((builder) => { builder .whereNull('locked_by') .orWhere('locked_at', '<', staleLockCutoffString); }) .orderBy([ 'ready_at', 'created_at', 'id' ]) .limit(limit); if (candidates.length === 0) { return { steps: [], nextStepReadyAt: await findNextPendingReadyAt(trx, staleLockCutoff) }; } const candidateIds = candidates.map(candidate => candidate.id); // 2. Try to lock those rows. await trx('automation_run_steps') .update({ locked_by: lockId, locked_at: nowString, started_at: nowString, updated_at: nowString }) .increment('step_attempts', 1) .whereIn('id', candidateIds) .where('status', 'pending') .where('ready_at', '<=', nowString) .where((builder) => { builder .whereNull('locked_by') .orWhere('locked_at', '<', staleLockCutoffString); }); // 3. Select any rows we successfully locked. const rows: StepToRunRow[] = await trx('automation_run_steps as step') .select( 'step.id as id', 'step.locked_by as locked_by', 'step.automation_run_id as automation_run_id', 'run.automation_id as automation_id', 'automation.slug as automation_slug', 'automation.status as automation_status', 'run.member_id as member_id', 'run.member_email as member_email', 'action.id as action_id', 'revision.id as automation_action_revision_id', 'action.type as type', 'step.ready_at as ready_at', 'step.step_attempts as step_attempts', 'revision.wait_hours as wait_hours', 'revision.email_subject as email_subject', 'revision.email_lexical as email_lexical', 'revision.email_design_setting_id as email_design_setting_id' ) .innerJoin('automation_runs as run', 'run.id', 'step.automation_run_id') .innerJoin('automations as automation', 'automation.id', 'run.automation_id') .innerJoin('automation_action_revisions as revision', 'revision.id', 'step.automation_action_revision_id') .innerJoin('automation_actions as action', 'action.id', 'revision.action_id') .whereIn('step.id', candidateIds) .where('step.locked_by', lockId) .orderBy([ 'step.ready_at', 'step.created_at', 'step.id' ]); return { steps: rows.map(row => buildStepToRun(row)), nextStepReadyAt: await findNextPendingReadyAt(trx, staleLockCutoff) }; } async function findNextPendingReadyAt(trx: Knex.Transaction, staleLockCutoff: Readonly): Promise { const row = await trx('automation_run_steps') .select({next_ready_at: 'ready_at'}) .where('status', 'pending') .where((builder) => { builder .whereNull('locked_by') .orWhere('locked_at', '<', toDatabaseDate(staleLockCutoff)); }) .orderBy('ready_at') .first(); return row?.next_ready_at ? fromDatabaseDate(row.next_ready_at) : null; } function buildStepToRun(row: ReadonlyDeep): AutomationStepToRun { const base = { id: row.id, step_attempts: row.step_attempts, ready_at: fromDatabaseDate(row.ready_at), locked_by: row.locked_by, automation_run_id: row.automation_run_id, automation_id: row.automation_id, automation_slug: row.automation_slug, automation_status: row.automation_status, member_id: row.member_id, member_email: row.member_email, action_id: row.action_id, automation_action_revision_id: row.automation_action_revision_id }; switch (row.type) { case 'wait': return { ...base, type: 'wait', wait_hours: requireValue(row, 'wait_hours') }; case 'send_email': return { ...base, type: 'send_email', email_subject: requireValue(row, 'email_subject'), email_lexical: requireValue(row, 'email_lexical'), email_design_setting_id: row.email_design_setting_id }; default: throw new errors.InternalServerError({ message: `Unexpected action type from database: ${row.type}` }); } } async function findFirstActionRevision(trx: Knex.Transaction, memberStatus: 'free' | 'paid'): Promise { const automationSlug: NonNullable = MEMBER_WELCOME_EMAIL_SLUGS[memberStatus]; const row = await trx('automations as automation') .select( 'automation.id as automation_id', 'actions.id as action_id', 'revisions.id as automation_action_revision_id', 'actions.type as type', 'revisions.wait_hours as wait_hours' ) .innerJoin('automation_actions as actions', 'actions.automation_id', 'automation.id') .innerJoin('automation_action_revisions as revisions', 'revisions.action_id', 'actions.id') .where('automation.slug', automationSlug) .where('automation.status', 'active') .whereNull('actions.deleted_at') .whereNotExists(trx('automation_action_edges as edge') .select('edge.target_action_id') .innerJoin('automation_actions as source_actions', 'source_actions.id', 'edge.source_action_id') .whereNull('source_actions.deleted_at') .where('edge.target_action_id', trx.ref('actions.id'))) .where('revisions.created_at', trx('automation_action_revisions') .max('created_at') .where('action_id', trx.ref('actions.id'))) .orderBy([ 'actions.created_at', 'actions.id' ]) .first(); return row ?? null; } async function finishStepAndEnqueueNext( trx: Knex.Transaction, options: Readonly<{ step: Pick; fakeWaitHoursMultiplier: number | null; }> ): Promise { const { step, fakeWaitHoursMultiplier } = options; const didFinish = await markStepTerminal(trx, step, 'finished'); if (!didFinish) { return null; } if (!await isRunAutomationActive(trx, step.automation_run_id)) { return null; } const next = await findNextActionRevision(trx, step.action_id); if (!next) { return null; } const now = new Date(); const nextReadyAt = getReadyAtForAction(next, now, fakeWaitHoursMultiplier); await insertRunStep(trx, { automationRunId: step.automation_run_id, automationActionRevisionId: next.automation_action_revision_id, now, readyAt: nextReadyAt }); return nextReadyAt; } async function findNextActionRevision(trx: Knex.Transaction, sourceActionId: string): Promise { const row = await trx('automation_action_edges as edge') .select( 'action.id as action_id', 'revision.id as automation_action_revision_id', 'action.type as type', 'revision.wait_hours as wait_hours' ) .innerJoin('automation_actions as action', 'action.id', 'edge.target_action_id') .innerJoin('automation_action_revisions as revision', 'revision.action_id', 'action.id') .where('edge.source_action_id', sourceActionId) .whereNull('action.deleted_at') .where('revision.created_at', trx('automation_action_revisions') .max('created_at') .where('action_id', trx.ref('action.id'))) .orderBy('revision.created_at', 'desc') .orderBy('revision.id', 'desc') .first(); return row ?? null; } async function markStepTerminal( trx: Knex.Transaction, step: Pick, status: AutomationStepTerminalStatus ): Promise { const nowString = toDatabaseDate(new Date()); return await updateStep(trx, step, { status, finished_at: nowString, updated_at: nowString }); } async function retryStep( trx: Knex.Transaction, step: Pick, retryAt: Readonly ): Promise { if (!await isStepRunAutomationActive(trx, step.id)) { await markStepTerminal(trx, step, 'automation disabled'); return false; } const nowString = toDatabaseDate(new Date()); return await updateStep(trx, step, { status: 'pending', started_at: null, finished_at: null, ready_at: toDatabaseDate(retryAt), updated_at: nowString }); } function getReadyAtForAction( action: ReadonlyDeep>, now: Readonly, fakeWaitHoursMultiplier: number | null ): Date { switch (action.type) { case 'wait': { const waitHours = requireValue({ ...action, id: action.action_id }, 'wait_hours'); const waitMs = waitHours * (fakeWaitHoursMultiplier || HOUR_MS); return new Date(now.getTime() + waitMs); } case 'send_email': return now; default: { const _exhaustive: never = action.type; throw new errors.IncorrectUsageError({ message: `Unexpected action type ${_exhaustive}` }); } } } async function cancelCancelablePendingStepsForAutomation( trx: Knex.Transaction, automationId: string, now: Readonly ): Promise { const nowString = toDatabaseDate(now); const staleLockCutoff = toDatabaseDate(getStaleLockCutoff(now)); await trx('automation_run_steps') .update({ status: 'automation disabled', finished_at: nowString, updated_at: nowString, locked_by: null, locked_at: null }) .where('status', 'pending') .whereIn('automation_run_id', trx('automation_runs') .select('id') .where('automation_id', automationId)) .where((builder) => { builder .whereNull('locked_by') .orWhere('locked_at', '<', staleLockCutoff); }); } async function isStepRunAutomationActive(trx: Knex.Transaction, stepId: string): Promise { const query = trx('automation_run_steps as step') .select(trx.raw('1')) .innerJoin('automation_runs as run', 'run.id', 'step.automation_run_id') .innerJoin('automations as automation', 'automation.id', 'run.automation_id') .where('step.id', stepId) .where('automation.status', 'active'); return await selectExists(trx, query); } async function isRunAutomationActive(trx: Knex.Transaction, automationRunId: string): Promise { const query = trx('automation_runs as run') .select(trx.raw('1')) .innerJoin('automations as automation', 'automation.id', 'run.automation_id') .where('run.id', automationRunId) .where('automation.status', 'active'); return await selectExists(trx, query); } async function selectExists(trx: Knex.Transaction, query: Knex.QueryBuilder): Promise { const row = await trx .select<{exists: boolean | number | string}>(trx.raw('exists ? as `exists`', [query])) .first(); return Boolean(Number(row?.exists)); } /** * Update a step. Returns whether the update succeeded. * * Should only update locked steps to avoid race conditions. Imagine the following scenario: * * 1. A step is locked by Worker A. * 2. The lock expires. * 3. The step is locked by Worker B. * 4. Worker A finishes its work. * * Worker A has lost its lock, so it shouldn't be updating the step any more. */ async function updateStep( trx: Knex.Transaction, step: Pick, attrs: { status: string; started_at?: string | null; finished_at?: string | null; ready_at?: string; updated_at: string; } ): Promise { /* eslint-disable camelcase */ const {started_at, finished_at, ready_at} = attrs; const changes = await trx('automation_run_steps') .update({ status: attrs.status, updated_at: attrs.updated_at, locked_by: null, locked_at: null, ...(started_at === undefined ? {} : {started_at}), ...(finished_at === undefined ? {} : {finished_at}), ...(ready_at === undefined ? {} : {ready_at}) }) .where('id', step.id) .where('status', 'pending') .where('locked_by', step.locked_by); /* eslint-enable camelcase */ return changes >= 1; } async function loadAutomation(trx: Knex.Transaction, automationId: string): Promise { const row = await trx('automations') .select('id', 'slug', 'name', 'status', 'created_at', 'updated_at') .where('id', automationId) .first(); return row ?? null; } async function loadAutomationBySlug(trx: Knex.Transaction, slug: string): Promise { const row = await trx('automations') .select('id', 'slug', 'name', 'status', 'created_at', 'updated_at') .where('slug', slug) .first(); return row ?? null; } async function loadAutomations(trx: Knex.Transaction): Promise { return await trx('automations') .select('id', 'slug', 'name', 'status', 'created_at', 'updated_at') .orderBy('name'); } async function updateAutomation(trx: Knex.Transaction, automation: AutomationRow): Promise { await trx('automations') .update({ status: automation.status, updated_at: automation.updated_at }) .where('id', automation.id); return requireAutomation(await loadAutomation(trx, automation.id), automation.id); } async function replaceAutomationGraph(trx: Knex.Transaction, automationId: string, submittedActions: AutomationAction[], edges: AutomationEdge[]): Promise { const existingActions = await loadAutomationActionRows(trx, automationId); const existingActionById = new Map(existingActions.map(action => [action.id, action])); const actions = (await resolveEmailDesignSettingIds(trx, submittedActions)).map(formatActionOnWrite); const submittedActionIds = new Set(actions.map(action => action.id)); const actionIdsWithOwners = await loadActionIdsWithOwners(trx, [...submittedActionIds]); const latestRevisionByActionId = new Map((await loadLatestActionRevisions(trx, [...submittedActionIds])).map(revision => [revision.action_id, revision])); const now = toDatabaseDate(new Date()); const actionsToInsert: ActionToInsert[] = []; const revisionsToInsert: ActionRevisionToInsert[] = []; for (const action of actions) { const existingAction = existingActionById.get(action.id); if (existingAction) { if (existingAction.type !== action.type) { throw new errors.ValidationError({ message: tpl(messages.conflictingAutomationActionType, { actionId: action.id }), property: 'actions.type' }); } } else { if (actionIdsWithOwners.has(action.id)) { throw new errors.ValidationError({ message: tpl(messages.conflictingAutomationActionId, { actionId: action.id }), property: 'actions.id' }); } actionsToInsert.push({ id: action.id, created_at: now, updated_at: now, automation_id: automationId, type: action.type }); } const latestRevision = latestRevisionByActionId.get(action.id); if (shouldInsertActionRevision(action, latestRevision)) { revisionsToInsert.push({ actionId: action.id, action, createdAt: getNextRevisionCreatedAt(latestRevision?.created_at ?? null, now) }); } } await insertActions(trx, actionsToInsert); await insertActionRevisions(trx, revisionsToInsert); const actionIdsToSoftDelete = existingActions .filter(existingAction => !submittedActionIds.has(existingAction.id)) .map(existingAction => existingAction.id); await softDeleteActions(trx, actionIdsToSoftDelete, now); await deleteAutomationEdges(trx, automationId); await insertActionEdges(trx, edges); } function formatActionOnWrite(action: AutomationAction): AutomationAction { if (action.type !== 'send_email') { return action; } return { ...action, data: { ...action.data, email_lexical: urlUtils.lexicalToTransformReady(action.data.email_lexical, { nodes: lexicalLib.nodes, transformMap: lexicalLib.urlTransformMap }) } }; } async function resolveEmailDesignSettingIds(trx: Knex.Transaction, actions: ReadonlyArray): Promise { if (!actions.some(action => action.type === 'send_email' && action.data.email_design_setting_id === DEFAULT_EMAIL_DESIGN_SETTING_REFERENCE)) { return [...actions]; } const defaultEmailDesignSettingId = await loadDefaultEmailDesignSettingId(trx); return actions.map((action) => { if (action.type !== 'send_email' || action.data.email_design_setting_id !== DEFAULT_EMAIL_DESIGN_SETTING_REFERENCE) { return action; } return { ...action, data: { ...action.data, email_design_setting_id: defaultEmailDesignSettingId } }; }); } async function loadDefaultEmailDesignSettingId(trx: Knex.Transaction): Promise { const row = await trx('email_design_settings') .select('id') .where('slug', DEFAULT_EMAIL_DESIGN_SETTING_SLUG) .first(); if (!row?.id) { throw new errors.InternalServerError({ message: tpl(messages.defaultEmailDesignSettingNotFound) }); } return row.id; } async function loadAutomationActionRows(trx: Knex.Transaction, automationId: string): Promise>> { return await trx('automation_actions') .select('id', 'type') .where('automation_id', automationId) .whereNull('deleted_at'); } async function loadActionIdsWithOwners(trx: Knex.Transaction, actionIds: ReadonlyArray): Promise> { if (actionIds.length === 0) { return new Set(); } const rows = await trx('automation_actions') .select('id') .whereIn('id', actionIds); return new Set(rows.map(row => row.id)); } async function insertActions(trx: Knex.Transaction, actions: ReadonlyArray): Promise { if (actions.length === 0) { return; } await trx('automation_actions').insert(actions); } function shouldInsertActionRevision(action: AutomationAction, latestRevision: ActionRevisionRow | undefined): boolean { if (!latestRevision) { return true; } return !dequal(buildRevisionActionData(action, latestRevision), action.data); } function buildRevisionActionData(action: AutomationAction, revision: ActionRevisionRow): ExclusifyUnion { switch (action.type) { case 'wait': return { wait_hours: revision.wait_hours }; case 'send_email': return { email_subject: revision.email_subject, email_lexical: revision.email_lexical, email_design_setting_id: revision.email_design_setting_id }; default: { const _exhaustive: never = action; throw new errors.InternalServerError({ message: `Unhandled action type: ${_exhaustive}` }); } } } async function loadLatestActionRevisions( trx: Knex.Transaction, actionIds: ReadonlyArray ): Promise { if (actionIds.length === 0) { return []; } const latestRevisionDates = trx('automation_action_revisions') .select('action_id') .max({created_at: 'created_at'}) .whereIn('action_id', actionIds) .groupBy('action_id') .as('latest_revision_dates'); return await trx('automation_action_revisions') .select( 'automation_action_revisions.action_id', 'automation_action_revisions.created_at', 'automation_action_revisions.wait_hours', 'automation_action_revisions.email_subject', 'automation_action_revisions.email_lexical', 'automation_action_revisions.email_design_setting_id' ) .innerJoin(latestRevisionDates, function () { this.on('automation_action_revisions.action_id', 'latest_revision_dates.action_id') .andOn('automation_action_revisions.created_at', 'latest_revision_dates.created_at'); }); } async function softDeleteActions( trx: Knex.Transaction, actionIds: ReadonlyArray, deletedAt: string ): Promise { if (actionIds.length === 0) { return; } await trx('automation_actions') .update({ deleted_at: deletedAt, updated_at: deletedAt }) .whereIn('id', actionIds); } async function insertActionRevisions( trx: Knex.Transaction, revisions: ReadonlyArray ): Promise { if (revisions.length === 0) { return; } await trx('automation_action_revisions').insert( revisions.map(({actionId, action, createdAt}) => buildActionRevision(actionId, action, createdAt)) ); } function getNextRevisionCreatedAt(latestCreatedAt: DatabaseDate | null, requestedCreatedAt: string) { if (!latestCreatedAt) { return toDatabaseDate(requestedCreatedAt); } const requestedTime = fromDatabaseDate(requestedCreatedAt).getTime(); const latestTime = fromDatabaseDate(latestCreatedAt).getTime(); if (requestedTime > latestTime) { return toDatabaseDate(requestedCreatedAt); } return toDatabaseDate(new Date(latestTime + 1000)); } function buildActionRevision(actionId: string, action: AutomationAction, createdAt: string) { switch (action.type) { case 'wait': return { id: ObjectId().toString(), created_at: createdAt, action_id: actionId, wait_hours: action.data.wait_hours, email_subject: null, email_lexical: null, email_design_setting_id: null }; case 'send_email': return { id: ObjectId().toString(), created_at: createdAt, action_id: actionId, wait_hours: null, email_subject: action.data.email_subject, email_lexical: action.data.email_lexical, email_design_setting_id: action.data.email_design_setting_id }; default: { const _exhaustive: never = action; throw new errors.InternalServerError({ message: `Unexpected action type ${_exhaustive}` }); } } } async function deleteAutomationEdges(trx: Knex.Transaction, automationId: string): Promise { await trx('automation_action_edges') .delete() .whereIn('source_action_id', trx('automation_actions') .select('id') .where('automation_id', automationId)); } async function insertActionEdges(trx: Knex.Transaction, edges: ReadonlyArray): Promise { if (edges.length === 0) { return; } await trx('automation_action_edges').insert(edges.map(edge => ({ source_action_id: edge.source_action_id, target_action_id: edge.target_action_id }))); } function requireAutomation(automation: AutomationRow | null, id: string): AutomationRow { if (!automation) { throw new errors.InternalServerError({ message: `Updated automation "${id}" could not be loaded.` }); } return automation; } async function buildAutomation(trx: Knex.Transaction, automation: AutomationRow): Promise { const actionRows = await loadActionRows(trx, automation.id); const actionStats = await loadActionStats(trx, actionRows.map(row => row.id)); const edgeRows = await loadEdgeRows(trx, automation.id); return { ...buildAutomationSummary(automation), actions: actionRows.map(row => buildActionPayload(row, actionStats.get(row.id) ?? null)), edges: edgeRows.map(row => buildEdgePayload(row)) }; } function buildAutomationSummary(automation: AutomationRow): AutomationSummary { return { id: automation.id, slug: automation.slug, name: automation.name, status: automation.status, created_at: serializeDate(automation.created_at), updated_at: serializeDate(automation.updated_at) }; } function serializeDate(date: DatabaseDate) { const normalizedDate = fromDatabaseDate(date); normalizedDate.setMilliseconds(0); return normalizedDate.toISOString(); } async function loadActionRows(trx: Knex.Transaction, automationId: string): Promise { return await trx('automation_actions as a') .select( 'a.id as id', 'a.type as type', 'r.wait_hours as wait_hours', 'r.email_subject as email_subject', 'r.email_lexical as email_lexical', 'r.email_design_setting_id as email_design_setting_id' ) .innerJoin('automation_action_revisions as r', 'r.action_id', 'a.id') .where('a.automation_id', automationId) .whereNull('a.deleted_at') .where('r.created_at', trx('automation_action_revisions') .max('created_at') .where('action_id', trx.ref('a.id'))) .orderBy([ 'a.created_at', 'a.id' ]); } async function loadActionStats( trx: Knex.Transaction, actionIds: ReadonlyArray ): Promise> { if (actionIds.length === 0) { return new Map(); } const rows: ActionStatsRow[] = await trx('automation_action_revisions') .select('action_id') .sum({ email_clicked_count: 'email_clicked_count', email_sent_count: 'email_sent_count', email_opened_count: 'email_opened_count' }) .whereIn('action_id', actionIds) .groupBy('action_id'); return new Map(rows.map(row => [row.action_id, buildEmailStats(row)])); } async function loadEdgeRows(trx: Knex.Transaction, automationId: string): Promise { return await trx('automation_action_edges as e') .select('e.source_action_id', 'e.target_action_id') .innerJoin('automation_actions as source_action', (join) => { join .on('source_action.id', 'e.source_action_id') .onNull('source_action.deleted_at'); }) .innerJoin('automation_actions as target_action', (join) => { join .on('target_action.id', 'e.target_action_id') .onNull('target_action.deleted_at') .on('target_action.automation_id', 'source_action.automation_id'); }) .where('source_action.automation_id', automationId) .orderBy([ 'e.source_action_id', 'e.target_action_id' ]); } function buildActionPayload(row: ActionRow, stats: AutomationEmailStats | null): AutomationAction { switch (row.type) { case 'wait': return { id: row.id, type: 'wait', data: { wait_hours: requireValue(row, 'wait_hours') } }; case 'send_email': return { id: row.id, type: 'send_email', data: { email_subject: requireValue(row, 'email_subject'), email_lexical: urlUtils.transformReadyToAbsolute(requireValue(row, 'email_lexical')), email_design_setting_id: requireValue(row, 'email_design_setting_id') }, stats: stats ?? EMPTY_EMAIL_STATS }; } } const EMPTY_EMAIL_STATS: AutomationEmailStats = { email_clicked_count: 0, email_sent_count: 0, email_opened_count: 0, opened_rate: null, clicked_rate: null }; function buildEmailStats(row: ActionStatsRow): AutomationEmailStats { const emailClickedCount = row.email_clicked_count ?? 0; const emailSentCount = row.email_sent_count ?? 0; const emailOpenedCount = row.email_opened_count ?? 0; return { email_clicked_count: emailClickedCount, email_sent_count: emailSentCount, email_opened_count: emailOpenedCount, opened_rate: emailSentCount ? Math.round(emailOpenedCount / emailSentCount * 100) : null, clicked_rate: emailSentCount ? Math.round(emailClickedCount / emailSentCount * 100) : null }; } function requireValue< RowT extends {id: string, type: string}, FieldT extends keyof RowT >( row: RowT, field: FieldT ): NonNullable { const value = row[field]; if ((value === null) || (value === undefined)) { throw new errors.InternalServerError({ message: tpl(messages.invalidAutomationActionRevision, { actionId: row.id, actionType: row.type, field }) }); } return value; } function buildEdgePayload(edge: EdgeRow): AutomationEdge { return { source_action_id: edge.source_action_id, target_action_id: edge.target_action_id }; } function buildPagination(total: number) { return { page: 1, pages: 1, limit: 'all' as const, total, prev: null, next: null }; }