diff --git a/.agents/skills/databuddy-internal/SKILL.md b/.agents/skills/databuddy-internal/SKILL.md index f9ee0f53eb..cfaa83462e 100644 --- a/.agents/skills/databuddy-internal/SKILL.md +++ b/.agents/skills/databuddy-internal/SKILL.md @@ -193,7 +193,7 @@ Read [codebase-map.md](./references/codebase-map.md) when you need deeper routin ### Database work -- Postgres schema: `packages/db/src/drizzle/schema.ts` +- Postgres schemas: `packages/db/src/drizzle/schema/` (`index.ts` barrel) - Relations: `packages/db/src/drizzle/relations.ts` - Drizzle client: `packages/db/src/client.ts` - Production `DATABASE_URL` may already target PgBouncer; inspect both the process pool and PgBouncer queues before attributing API timeouts to PostgreSQL. diff --git a/.agents/skills/databuddy-internal/references/codebase-map.md b/.agents/skills/databuddy-internal/references/codebase-map.md index b350f33603..8de7ddab6a 100644 --- a/.agents/skills/databuddy-internal/references/codebase-map.md +++ b/.agents/skills/databuddy-internal/references/codebase-map.md @@ -73,7 +73,7 @@ Use this file when the task spans multiple packages or when the right edit locat - Postgres schema and relations - ClickHouse client and schema - Key files: - - [`packages/db/src/drizzle/schema.ts`](/Users/iza/Dev/Databuddy/packages/db/src/drizzle/schema.ts) + - [`packages/db/src/drizzle/schema/index.ts`](/Users/iza/Dev/Databuddy/packages/db/src/drizzle/schema/index.ts) - [`packages/db/src/drizzle/relations.ts`](/Users/iza/Dev/Databuddy/packages/db/src/drizzle/relations.ts) - [`packages/db/src/client.ts`](/Users/iza/Dev/Databuddy/packages/db/src/client.ts) — strips `sslrootcert=system` from `DATABASE_URL` before `pg` Pool: libpq uses it for the OS trust store, but node-postgres treats `sslrootcert` as a file path and throws `ENOENT` on path `"system"`. - [`packages/db/src/clickhouse/client.ts`](/Users/iza/Dev/Databuddy/packages/db/src/clickhouse/client.ts) diff --git a/apps/dashboard/app/(main)/organizations/components/business-context-editor.tsx b/apps/dashboard/app/(main)/organizations/components/business-context-editor.tsx index b0c94417a2..4c11f34919 100644 --- a/apps/dashboard/app/(main)/organizations/components/business-context-editor.tsx +++ b/apps/dashboard/app/(main)/organizations/components/business-context-editor.tsx @@ -10,6 +10,8 @@ import { type BusinessContextSettings, businessContextIsGenerating, formatBusinessTeamContext, + formatBusinessMeasurementPlans, + businessMeasurementPlansSchema, } from "@databuddy/shared/organization-business-context"; import { Button, Card, Field, Textarea, dayjs } from "@databuddy/ui"; import { Accordion, Dialog, DropdownMenu } from "@databuddy/ui/client"; @@ -25,6 +27,7 @@ import { useEffect, useRef, useState } from "react"; import { TopBar } from "@/components/layout/top-bar"; import { getUserFacingErrorMessage } from "@/lib/user-facing-error"; import { useBusinessContextDraft } from "./use-business-context-draft"; +import { MeasurementPlanEditor } from "./measurement-plan-editor"; const emptyTeamContext: BusinessTeamContext = { priority: "", @@ -190,6 +193,8 @@ export function BusinessContextEditor({ const content = draft?.content ?? profile?.content ?? ""; const teamContext = draft?.teamContext ?? profile?.teamContext ?? emptyTeamContext; + const measurementPlans = + draft?.measurementPlans ?? profile?.measurementPlans ?? []; const generationWebsite = websites.find( (site) => site.id === generation?.websiteId && site.domain === generation.domain @@ -207,7 +212,9 @@ export function BusinessContextEditor({ (content.trim() !== (profile?.content ?? "") || Boolean(draftGeneration) || formatBusinessTeamContext(teamContext) !== - formatBusinessTeamContext(profile?.teamContext)); + formatBusinessTeamContext(profile?.teamContext) || + JSON.stringify(measurementPlans) !== + JSON.stringify(profile?.measurementPlans ?? [])); const conflict = dirty && draft.revision !== revision; const activeGeneration = businessContextIsGenerating(settings); const generating = isRequesting || activeGeneration; @@ -233,11 +240,20 @@ export function BusinessContextEditor({ const teamTooLong = Object.values(teamContext).some( (value) => value.trim().length > BUSINESS_CONTEXT_TEAM_FIELD_LIMIT ); + const plansValid = + businessMeasurementPlansSchema.safeParse(measurementPlans).success; + const bindingsValid = measurementPlans.every((plan) => + websites.some( + (site) => site.id === plan.websiteId && site.domain === plan.domain + ) + ); const saveDisabled = !(ready && canEdit && dirty) || conflict || tooLong || teamTooLong || + !plansValid || + !bindingsValid || isSaving || review !== null; const reviewedProfile = review?.kind === "history" ? review.profile : profile; @@ -247,6 +263,10 @@ export function BusinessContextEditor({ : (reviewedProfile?.content ?? ""); const reviewTeam = review?.kind === "generation" ? teamContext : reviewedProfile?.teamContext; + const reviewPlans = + review?.kind === "generation" + ? measurementPlans + : reviewedProfile?.measurementPlans; useEffect(() => { if ( @@ -265,6 +285,7 @@ export function BusinessContextEditor({ revision, generationId: readyGeneration.id, teamContext: profile?.teamContext, + measurementPlans: profile?.measurementPlans, }); }, [ ready, @@ -274,6 +295,7 @@ export function BusinessContextEditor({ readyGeneration, revision, profile?.teamContext, + profile?.measurementPlans, setDraft, ]); @@ -333,6 +355,7 @@ export function BusinessContextEditor({ content: content.trim(), revision: draft.revision, teamContext, + measurementPlans, ...(draftGeneration ? { generationId: draftGeneration.id } : {}), }), "Changes saved" @@ -676,6 +699,30 @@ export function BusinessContextEditor({ )} + { + setDraft({ + ...(draft ?? { content, revision, teamContext }), + measurementPlans: plans, + }); + setNotice(""); + }} + /> + {dirty && !plansValid && ( +

+ Complete the outcome name and both event names before saving. Event + names and namespace can contain up to 256 characters. +

+ )} + {dirty && !bindingsValid && ( +

+ Update or remove definitions for changed or unavailable websites + before saving. +

+ )} @@ -706,16 +753,24 @@ export function BusinessContextEditor({ {review?.kind === "generation" ? "Using this draft replaces your local text. You can edit it before saving." : review?.kind === "history" - ? "Restoring replaces the saved brief and your current edits. Your current saved version stays in history." + ? "Restoring replaces the saved brief, team context, event definitions, and your current edits. Your current saved version stays in history." : "Your edits are still in the editor. Choose which version to keep working on."} @@ -788,6 +843,7 @@ export function BusinessContextEditor({ revision, generationId: pendingDraft.id, teamContext, + measurementPlans, }); setReview(null); editorRef.current?.focus(); diff --git a/apps/dashboard/app/(main)/organizations/components/measurement-plan-editor.tsx b/apps/dashboard/app/(main)/organizations/components/measurement-plan-editor.tsx new file mode 100644 index 0000000000..9aaa747514 --- /dev/null +++ b/apps/dashboard/app/(main)/organizations/components/measurement-plan-editor.tsx @@ -0,0 +1,311 @@ +"use client"; + +import type { + BusinessContextSettings, + BusinessMeasurementPlan, +} from "@databuddy/shared/organization-business-context"; +import { Button, Card, Field, Input } from "@databuddy/ui"; +import { Accordion, DropdownMenu } from "@databuddy/ui/client"; +import { CaretDownIcon } from "@databuddy/ui/icons"; +import { useState } from "react"; +import { AutocompleteInput } from "@/components/ui/autocomplete-input"; +import { useAutocompleteData } from "@/hooks/use-autocomplete"; + +interface MeasurementPlanEditorProps { + disabled: boolean; + onChange: (plans: BusinessMeasurementPlan[]) => void; + plans: BusinessMeasurementPlan[]; + websites: BusinessContextSettings["websites"]; +} + +const eventFields = [ + { key: "activationEvent", label: "Activation event" }, + { key: "returnEvent", label: "Return event" }, +] as const; + +export function MeasurementPlanEditor({ + websites, + plans, + disabled, + onChange, +}: MeasurementPlanEditorProps) { + const [websiteId, setWebsiteId] = useState(""); + const website = websites.find((site) => site.id === websiteId) ?? websites[0]; + const plan = plans.find((item) => item.websiteId === website?.id); + const catalog = useAutocompleteData(website?.id ?? "", !disabled && !!plan); + const events = catalog.data?.customEvents ?? []; + const domainMismatch = plan && website && plan.domain !== website.domain; + const update = ( + changes: Partial> + ) => { + if (disabled || !plan) { + return; + } + onChange( + plans.map((item) => + item.websiteId === plan.websiteId ? { ...item, ...changes } : item + ) + ); + }; + const toggleDefinition = () => { + if (disabled || !website) { + return; + } + onChange( + plan + ? plans.filter((item) => item.websiteId !== website.id) + : [ + ...plans, + { + websiteId: website.id, + domain: website.domain, + name: "", + activationEvent: "", + returnEvent: "", + horizonDays: 7, + }, + ] + ); + }; + + return ( + + + Activation and return + + Choose the events that mean someone got value and came back. Saved + definitions guide automatic investigations. Only identified profiles + can be measured. + + + + {plans + .filter( + (item) => !websites.some((site) => site.id === item.websiteId) + ) + .map((item) => ( +
+

+ {item.name || item.domain}: website unavailable. This definition + is inactive. +

+ {!disabled && ( + + )} +
+ ))} + {disabled ? ( + plans.length ? ( + plans.map((item) => { + const site = websites.find( + (candidate) => candidate.id === item.websiteId + ); + return ( +
+

+ {item.name || "Unnamed outcome"} +

+

{item.domain}

+ {site && site.domain !== item.domain && ( +

+ Website domain changed to {site.domain}. This definition + is inactive until updated. +

+ )} +

+ Activation: {item.activationEvent || "Not set"} +

+

+ Return: {item.returnEvent || "Not set"} within{" "} + {item.horizonDays} days +

+ {item.namespace && ( +

+ Namespace: {item.namespace} +

+ )} +
+ ); + }) + ) : ( +

+ No definitions configured. +

+ ) + ) : website ? ( + <> +
+ {websites.length > 1 ? ( + + + } + > + {website.domain} + + + + + {websites.map((site) => ( + + {site.domain} + + ))} + + + + ) : ( +

+ {website.domain} +

+ )} + +
+ {plan ? ( +
+ {domainMismatch && ( +
+

+ This definition is bound to {plan.domain}. Update it to{" "} + {website.domain} before saving. +

+ +
+ )} + + Business outcome + update({ name: event.target.value })} + placeholder="Name this outcome for your team" + value={plan.name} + /> + +
+ {eventFields.map(({ key, label }) => ( + + {label} + update({ [key]: value })} + placeholder="Exact event name" + suggestions={events} + value={plan[key]} + /> + + {catalog.isError + ? "Catalog unavailable; enter an exact name." + : catalog.isPending + ? "Loading event names; you can keep typing." + : plan[key] + ? events.includes(plan[key]) + ? "Seen in the recent event catalog." + : "Not seen recently. Check that this event is recorded." + : "Choose a recent event or type an exact name."} + + + ))} +
+ + } + > + Return within {plan.horizonDays} days + + + + + update({ horizonDays: value === "30" ? 30 : 7 }) + } + value={String(plan.horizonDays)} + > + + 7 days + + + 30 days + + + + + + + Advanced{plan.namespace ? " · Namespace set" : ""} + + + + Namespace (optional) + + update({ namespace: event.target.value || undefined }) + } + placeholder="Exact namespace" + spellCheck={false} + value={plan.namespace ?? ""} + /> + + + +
+ ) : ( +

+ {plans.length >= 20 + ? "Up to 20 website definitions are supported." + : "No definition for this website. Add one to choose the outcome and events."} +

+ )} + + ) : ( +

+ Add a website to define activation and return. +

+ )} +
+
+ ); +} diff --git a/apps/dashboard/app/(main)/organizations/components/use-business-context-draft.ts b/apps/dashboard/app/(main)/organizations/components/use-business-context-draft.ts index a05e868e70..a9580ecc9a 100644 --- a/apps/dashboard/app/(main)/organizations/components/use-business-context-draft.ts +++ b/apps/dashboard/app/(main)/organizations/components/use-business-context-draft.ts @@ -2,6 +2,7 @@ import { businessContextEditSchema, + businessMeasurementPlanSchema, type BusinessContextEdit, } from "@databuddy/shared/organization-business-context"; import { useCallback, useEffect, useState } from "react"; @@ -10,6 +11,17 @@ import { z } from "zod"; // Keep invalid/unfinished input recoverable too; saving applies the real limits. const recoverySchema = businessContextEditSchema.extend({ content: z.string().max(100_000), + measurementPlans: z + .array( + businessMeasurementPlanSchema.extend({ + name: z.string().max(1000), + activationEvent: z.string().max(1000), + returnEvent: z.string().max(1000), + namespace: z.string().max(1000).optional(), + }) + ) + .max(20) + .optional(), teamContext: z .object({ priority: z.string().max(10_000), diff --git a/apps/dashboard/test/e2e/specs/regressions/measurement-plan.spec.ts b/apps/dashboard/test/e2e/specs/regressions/measurement-plan.spec.ts new file mode 100644 index 0000000000..eaf0ab3b3c --- /dev/null +++ b/apps/dashboard/test/e2e/specs/regressions/measurement-plan.spec.ts @@ -0,0 +1,67 @@ +import { expect, test } from "@/test/e2e/fixtures"; + +test("saves activation definitions through oRPC, recovers unfinished edits, and restores history", { + tag: "@regression", +}, async ({ authenticatedPage: page }) => { + await page.goto("/organizations/settings/business-context"); + await page.getByRole("button", { name: /^Add definition for/ }).click(); + const outcome = page.getByRole("textbox", { + name: "Business outcome", + exact: true, + }); + const activation = page.getByRole("combobox", { + name: "Activation event", + exact: true, + }); + const returning = page.getByRole("combobox", { + name: "Return event", + exact: true, + }); + await outcome.fill("Reports shared again"); + await expect( + page.getByRole("button", { name: "Save changes", exact: true }) + ).toBeDisabled(); + await page.reload(); + await expect(outcome).toHaveValue("Reports shared again"); + await activation.fill("report_shared"); + await returning.fill("report_opened"); + await returning.press("Escape"); + const saved = page.waitForResponse((response) => + response.url().endsWith("/rpc/businessContext/save") + ); + await page.getByRole("button", { name: "Save changes", exact: true }).click(); + expect((await saved).ok()).toBe(true); + await expect(page.getByText("Changes saved", { exact: true })).toBeVisible(); + await page.reload(); + await expect(activation).toHaveValue("report_shared"); + await expect(returning).toHaveValue("report_opened"); + await page.getByRole("button", { name: "Return window: 7 days" }).click(); + await page + .getByRole("menuitemradio", { name: "30 days", exact: true }) + .click(); + await page.getByRole("button", { name: "Save changes", exact: true }).click(); + await expect(page.getByText("Changes saved", { exact: true })).toBeVisible(); + await page.getByRole("button", { name: "History", exact: true }).click(); + await page + .getByRole("menuitem") + .filter({ hasText: /^Version/ }) + .first() + .click(); + await expect(page.getByRole("dialog").locator("ins")).toContainText("7"); + await page + .getByRole("button", { name: "Restore this version", exact: true }) + .click(); + await expect( + page.getByText("Version restored", { exact: true }) + ).toBeVisible(); + await expect( + page.getByRole("button", { name: "Return window: 7 days" }) + ).toBeVisible(); + await page.getByRole("button", { name: /^Remove definition for/ }).click(); + await page.getByRole("button", { name: "Save changes", exact: true }).click(); + await expect(page.getByText("Changes saved", { exact: true })).toBeVisible(); + await page.reload(); + await expect( + page.getByRole("button", { name: /^Add definition for/ }) + ).toBeVisible(); +}); diff --git a/apps/insights/src/business-aware-selection.test.ts b/apps/insights/src/business-aware-selection.test.ts index 851afac82e..5b4cf85355 100644 --- a/apps/insights/src/business-aware-selection.test.ts +++ b/apps/insights/src/business-aware-selection.test.ts @@ -634,3 +634,73 @@ describe("business-aware investigation selection", () => { ).toEqual(retry); }); }); + + +describe("saved activation measurement selection", () => { + const retention: DetectedSignal = { + ...outcome, + metric: "identified_retention", + subjectKey: "retention:synthetic", + label: "Reports shared again", + evidence: [ + "Native complete cohorts: 160/200 returned before, 80/200 after.", + ], + }; + it("avoids a selection call while preserving critical reliability and due work", async () => { + let calls = 0; + for (const dueSignalKey of [undefined, "goal:report-delivery"]) { + const selected = await planInvestigationsWithBusinessContext( + input, + [traffic, retention, outcome, error], + { + loadBusinessProfile: async () => ({ ...context, sources: [] }), + selectCandidates: async () => { + calls++; + throw new Error("Unexpected selection"); + }, + }, + false, + scope, + { reason: "manual", dueSignalKey } + ); + const keys = selected.map((candidate) => candidate.signal.signalKey); + expect(keys).toContain(retention.subjectKey!); + expect(keys).toContain(error.subjectKey!); + if (dueSignalKey) expect(keys[0]).toBe(dueSignalKey); + expect(keys).not.toContain("visitors"); + } + expect(calls).toBe(0); + }); + it.each([ + "team_reply", + "organization_profile", + ] as const)("allows %s context to supersede the saved measurement priority", async (kind) => { + let calls = 0; + const selected = await planInvestigationsWithBusinessContext( + input, + [traffic, retention, outcome], + { + loadBusinessProfile: async () => ({ + ...context, + sources: context.sources.map((source) => ({ ...source, kind })), + }), + selectCandidates: (params) => { + calls++; + return chooseInvestigationSignals( + params, + new MockLanguageModelV3({ + doGenerate: async () => response({ selections: [choice] }), + }) + ); + }, + }, + false, + scope, + { reason: "scheduled" } + ); + expect(calls).toBe(1); + expect(selected.map((candidate) => candidate.signal.signalKey)).toEqual([ + choice.signalKey, + ]); + }); +}); diff --git a/apps/insights/src/business-context.ts b/apps/insights/src/business-context.ts index 3e0a30f8ce..9cd48accfb 100644 --- a/apps/insights/src/business-context.ts +++ b/apps/insights/src/business-context.ts @@ -358,7 +358,8 @@ export async function loadWebsiteBusinessProfile( organizationProfileContext( value.profile, input.scope.organizationId, - input.allowRefresh ? new Date() : input.asOf + input.allowRefresh ? new Date() : input.asOf, + input.scope ) ) .catch((error) => @@ -380,10 +381,32 @@ export async function loadWebsiteBusinessProfile( export function organizationProfileContext( profile: OrganizationBusinessProfile | null, organizationId: string, - asOf: Date + asOf: Date, + scope?: Pick ): BusinessContext { const sources: BusinessSource[] = []; if (profile && Date.parse(profile.updatedAt) <= asOf.getTime()) { + const plan = profile.measurementPlans?.find( + (item) => + item.websiteId === scope?.websiteId && item.domain === scope.domain + ); + if (plan) { + const content = `Saved team-defined activation and return measurement (not emitter-code verification): ${JSON.stringify(plan)}. Native query: identified_profile_retention.`; + for (let offset = 0; offset < content.length; offset += 4000) { + sources.push({ + id: `organization-measurement-plan:${organizationId}:${plan.websiteId}:${offset / 4000}`, + kind: "organization_profile", + content: content.slice(offset, offset + 4000), + observedAt: profile.updatedAt, + author: "Team measurement definition", + origin: "team", + profileVersion: { + revision: profile.revision, + updatedAt: profile.updatedAt, + }, + }); + } + } const teamContext = formatBusinessTeamContext(profile.teamContext); for (let offset = 0; offset < teamContext.length; offset += 4000) { sources.push({ diff --git a/apps/insights/src/detection.ts b/apps/insights/src/detection.ts index f88113c32c..3fe14d2b84 100644 --- a/apps/insights/src/detection.ts +++ b/apps/insights/src/detection.ts @@ -3,6 +3,7 @@ import { normalizeCurrencyCode } from "@databuddy/shared/currency"; import type { InvestigationSignal, MatchedErrorContinuationMeasurement, + WeekOverWeekPeriod, } from "@databuddy/shared/insights"; import dayjs from "dayjs"; import timezonePlugin from "dayjs/plugin/timezone"; @@ -30,10 +31,12 @@ export interface DetectedSignal { direction: "up" | "down"; entityId?: string; entityLabel?: string; + evidence?: string[]; investigationObjective?: string; label: string; method: "behavior" | "zscore" | "wow"; metric: string; + period?: WeekOverWeekPeriod; severity: "critical" | "warning" | "info"; subjectKey?: string; } diff --git a/apps/insights/src/funnel-detection.ts b/apps/insights/src/funnel-detection.ts index 1953daf070..5823abe762 100644 --- a/apps/insights/src/funnel-detection.ts +++ b/apps/insights/src/funnel-detection.ts @@ -313,7 +313,7 @@ export function defaultFunnelGoalDeps( }; } -async function raceWithAbort( +export async function raceWithAbort( work: () => Promise, signal: AbortSignal ): Promise { @@ -321,7 +321,7 @@ async function raceWithAbort( let removeAbortListener: (() => void) | undefined; const stopped = new Promise((_resolve, reject) => { const onAbort = () => { - reject(signal.reason ?? new Error("Goal and funnel detection aborted")); + reject(signal.reason ?? new Error("Analytics detection aborted")); }; if (signal.aborted) { onAbort(); diff --git a/apps/insights/src/generation.ts b/apps/insights/src/generation.ts index 356184d42d..ec884fd37f 100644 --- a/apps/insights/src/generation.ts +++ b/apps/insights/src/generation.ts @@ -37,6 +37,7 @@ import { detectSignals, remeasureMetricSignal, } from "./detection"; +import { detectRetentionSignals } from "./measurement-plan"; import { detectFunnelGoalSignals, type FunnelGoalDeps, @@ -307,6 +308,7 @@ interface InvestigationRuntime { export interface InvestigationSources { detectDefinitionSignals: typeof detectFunnelGoalSignals; detectMetricSignals: typeof detectSignals; + detectRetentionSignals?: typeof detectRetentionSignals; detectRouteHealthSignals: typeof detectRouteHealthSignals; fetchAnnotations: ( websiteId: string, @@ -348,10 +350,20 @@ export function remeasureStoredSignal( abortSignal?: AbortSignal, dependencies: { funnelGoal?: FunnelGoalDeps; + retention?: Parameters[3]; query?: Parameters[2]; routeHealth?: RouteHealthDetectionDeps; } = {} ): Promise { + if (prior.signalKey.startsWith("retention:")) { + return detectRetentionSignals( + params, + today, + abortSignal, + dependencies.retention, + prior + ).then((signals) => signals[0] ?? null); + } return prior.signalKey.startsWith("goal:") || prior.signalKey.startsWith("funnel:") ? remeasureFunnelGoalSignal( @@ -466,6 +478,7 @@ export async function refreshInvestigationSignal(params: { } const productionInvestigationSources: InvestigationSources = { + detectRetentionSignals, loadBusinessProfile: loadWebsiteBusinessProfile, recallBusinessContext: recallWebsiteBusinessContext, detectDefinitionSignals: detectFunnelGoalSignals, @@ -605,6 +618,15 @@ async function discoverWebsiteSignals( sourceAbortSignal ) ), + detectSource( + "retention", + () => + runtime.sources.detectRetentionSignals?.( + detectParams, + asOf, + sourceAbortSignal + ) ?? Promise.resolve([]) + ), ] as const; const settledDetections = await Promise.allSettled(detectionTasks); const failedDetection = settledDetections.find( @@ -613,8 +635,13 @@ async function discoverWebsiteSignals( if (failedDetection?.status === "rejected") { throw discoveryController.signal.reason ?? failedDetection.reason; } - const [remeasuredDue, metricSignals, funnelGoalSignals, routeHealthSignals] = - await Promise.all(detectionTasks); + const [ + remeasuredDue, + metricSignals, + funnelGoalSignals, + routeHealthSignals, + retentionSignals, + ] = await Promise.all(detectionTasks); if ( due && remeasuredDue && @@ -636,6 +663,7 @@ async function discoverWebsiteSignals( ...metricSignals, ...funnelGoalSignals, ...routeHealthSignals, + ...retentionSignals, ]) { const key = signalKeyForDetectedSignal(signal); if (!signalsByKey.has(key)) { @@ -997,6 +1025,24 @@ export async function planInvestigationsWithBusinessContext( : disabled; // The shared profile already contains bounded, scoped PostgreSQL team replies. // Only selected subjects incur recall, analytics enrichment and investigation loops. + // Descriptive context, priorities, exclusions or replies can change what matters. + // Only a standalone saved measurement can skip contextual selection safely. + const plannedKeys = profile.sources.every((source) => + source.id.startsWith("organization-measurement-plan:") + ) + ? signals + .filter((signal) => signal.metric === "identified_retention") + .map(signalKeyForDetectedSignal) + : []; + + if (plannedKeys.length) { + // A saved exact measurement already supplies the question; preserve critical + // reliability and due work without spending a model call to rediscover it. + candidates = planCoveragePortfolio(signals, { + ...options, + selectedSignalKeys: plannedKeys, + }).map(toPlannedCandidate); + } const protectedCount = candidates.filter( (candidate) => candidate.signal.signalKey === options.dueSignalKey || @@ -1007,6 +1053,7 @@ export async function planInvestigationsWithBusinessContext( ).length; if ( sources.selectCandidates && + plannedKeys.length === 0 && profile.sources.length > 0 && (profile.status === "ready" || profile.status === "partial") && signals.length > 1 && diff --git a/apps/insights/src/investigation-flow.test.ts b/apps/insights/src/investigation-flow.test.ts index 187ed826bb..3b183a1559 100644 --- a/apps/insights/src/investigation-flow.test.ts +++ b/apps/insights/src/investigation-flow.test.ts @@ -3870,3 +3870,81 @@ describe("validateNumericGrounding", () => { ).not.toThrow(); }); }); + + +describe("identified-profile cohort publication", () => { + const comparison = + "Eligible identified profiles returning within seven days fell from 140/200 (70%) to 60/200 (30%)."; + const finish = { + title: "Report reuse fell", + summary: "Fewer identified profiles returned after sharing a report.", + rootCause: null, + evidence: [comparison], + evidenceRefs: [{ source: "provided", index: 0 }], + publish: true, + findingKind: "product_outcome", + publicationBasis: "measured_impact", + next: { + type: "resolve", + reason: "The measured change is useful; its cause remains unknown.", + }, + }; + const cohort: InvestigationSignal = { + ...signal, + signalKey: "retention:synthetic", + entity: { type: "cohort", id: "synthetic", label: "Shared reports" }, + metric: { + label: "Return within seven days", + format: "percent", + current: 30, + previous: 70, + }, + changePercent: -57.14, + }; + it("publishes a known-purpose cohort finding without a redundant data read or invented cause", async () => { + const model = outputModel(finish); + const result = await runInsightAgent( + { + appContext: appContext(), + signal: cohort, + evidence: [ + comparison, + "The team defines sharing a report as initial value and opening it later as reuse.", + ], + history: [], + otherOpenWork: [], + githubRepository: null, + }, + { model, tools: {} } + ); + expect(result.outcome).toMatchObject({ + publish: true, + findingKind: "product_outcome", + rootCause: null, + next: { type: "resolve" }, + }); + expect(model.doGenerateCalls).toHaveLength(1); + expect(result.toolCallCount).toBe(0); + }); + it("still rejects relabeling raw website traffic as a product loss", async () => { + await expect( + runInsightAgent( + { + appContext: appContext(), + signal: { + ...cohort, + signalKey: "visitors", + entity: { type: "website", id: "website", label: "Visitors" }, + }, + evidence: [comparison], + history: [], + otherOpenWork: [], + githubRepository: null, + }, + { model: outputModel(finish), tools: {} } + ) + ).rejects.toThrow( + "A website traffic signal is not a verified product loss" + ); + }); +}); diff --git a/apps/insights/src/investigation.ts b/apps/insights/src/investigation.ts index 90ecec6cb5..eb0340db1d 100644 --- a/apps/insights/src/investigation.ts +++ b/apps/insights/src/investigation.ts @@ -68,6 +68,7 @@ function metricFormat(metric: string): InsightMetric["format"] { if ( metric === "bounce_rate" || metric === "attribution_rate" || + metric === "identified_retention" || metric.startsWith("funnel:") || metric.startsWith("goal:") ) { @@ -117,6 +118,7 @@ function isDirectSignal(signal: DetectedSignal): boolean { signal.metric === "revenue" || signal.metric === "refund_amount" || signal.metric === "attribution_rate" || + signal.metric === "identified_retention" || signal.subjectKey?.includes(":referrer:") === true || signal.metric === "error_count" || signal.metric === "custom_event_count" || @@ -165,6 +167,7 @@ export function isInvestigationCandidate(signal: DetectedSignal): boolean { "product_revenue", "refund_amount", "attribution_rate", + "identified_retention", ].includes(signal.metric) || (isConversionDefinitionSignal(signal) && signal.current - signal.baseline >= 10 && @@ -206,6 +209,14 @@ export function rankSignals(signals: DetectedSignal[]): DetectedSignal[] { } function signalWindow(signal: DetectedSignal, lookbackDays: number) { + if (signal.period) { + return { + currentFrom: signal.period.current.from, + currentTo: signal.period.current.to, + previousFrom: signal.period.previous.from, + previousTo: signal.period.previous.to, + }; + } const detectedDay = dayjs(signal.detectedAt); if (signal.method === "zscore") { const baselineDates = signal.baselineDates ?? []; @@ -236,6 +247,13 @@ function entity(signal: DetectedSignal): InvestigationSignal["entity"] { const exactId = idParts.join(":"); const rawId = exactId.trim(); const id = boundedKey(rawId); + if (prefix === "retention" && signal.metric === "identified_retention") { + return { + type: "cohort", + id, + label: (signal.entityLabel ?? signal.label).slice(0, 120), + }; + } if (prefix === "funnel" && idParts.at(1) === "step") { return { type: "funnel_step", @@ -346,7 +364,7 @@ export function prepareInvestigation( ? { cohortMeasurement: candidate.cohortMeasurement } : {}), }; - const evidence: string[] = []; + const evidence: string[] = [...(candidate.evidence ?? [])]; if (candidate.definitionEvidence) { evidence.push(evidenceSummary(candidate.definitionEvidence)); } diff --git a/apps/insights/src/measurement-plan.test.ts b/apps/insights/src/measurement-plan.test.ts new file mode 100644 index 0000000000..7680df0695 --- /dev/null +++ b/apps/insights/src/measurement-plan.test.ts @@ -0,0 +1,237 @@ +import "@databuddy/test/env"; +import { describe, expect, it } from "bun:test"; +import type { executeQuery } from "@databuddy/ai/query"; +import type { BusinessMeasurementPlan } from "@databuddy/shared/organization-business-context"; +import dayjs from "dayjs"; +import { prepareInvestigation } from "./investigation"; +import { parseFrozenInvestigationPlan } from "./run-candidate-plan"; +import { organizationProfileContext } from "./business-context"; +import { + detectRetentionSignals, + measureActivationRetention, + measurementPlanKey, +} from "./measurement-plan"; + +const plan: BusinessMeasurementPlan = { + websiteId: "synthetic-site", + domain: "example.com", + name: "Shared reports", + activationEvent: "report_shared", + returnEvent: "report_opened", + horizonDays: 7, +}; +const asOf = dayjs("2026-09-09T12:00:00Z"); +const params = { websiteId: plan.websiteId, timezone: "UTC", lookbackDays: 7 }; + +function fixture( + options: { + eligible?: number; + before?: number; + after?: number; + incomplete?: number; + identity?: number; + } = {} +) { + const eligible = options.eligible ?? 200; + const before = options.before ?? 160; + const after = options.after ?? 80; + const incomplete = options.incomplete ?? 0; + const events = Math.ceil((eligible + incomplete) / (options.identity ?? 1)); + const query: typeof executeQuery = async (request) => { + const retained = request.from === "2026-08-18" ? before : after; + const row = { + cohort_from: request.from, + cohort_to: request.to, + observation_end: "2026-09-08", + cohort_start: dayjs.tz(request.from, "UTC").toISOString(), + cohort_end: dayjs.tz(request.to, "UTC").add(1, "day").toISOString(), + observed_before: "2026-09-09T00:00:00.000Z", + timezone: "UTC", + horizon_days: 7, + identity_basis: "direct_profile_id", + activation_basis: "first_in_cohort_window", + activated_profiles: eligible + incomplete, + eligible_profiles: eligible, + retained_profiles: retained, + not_retained_profiles: eligible - retained, + incomplete_profiles: incomplete, + activation_events: events, + identified_activation_events: eligible + incomplete, + unidentified_activation_events: events - eligible - incomplete, + }; + return [ + { ...row, row_type: "overall", cohort_date: null }, + { ...row, row_type: "cohort", cohort_date: request.from }, + ]; + }; + return query; +} + +describe("saved activation and return measurement", () => { + it("measures two independent complete cohorts in parallel native queries and preserves exact evidence", async () => { + let calls = 0; + const query: typeof executeQuery = async (...args) => { + calls++; + expect(args[0].type).toBe("identified_profile_retention"); + expect(args[0].projectId).toBe(plan.websiteId); + expect(args[0].filters).toContainEqual({ + field: "namespace", + op: "eq", + value: "production", + }); + return await fixture()(...args); + }; + const signals = await detectRetentionSignals(params, asOf, undefined, { + readPlan: async () => ({ ...plan, namespace: "production" }), + query, + }); + expect(calls).toBe(2); + expect(signals).toHaveLength(1); + expect(signals[0]).toMatchObject({ + current: 40, + baseline: 80, + metric: "identified_retention", + direction: "down", + }); + const prepared = prepareInvestigation(signals[0], 7); + expect(prepared.signal.period).toEqual({ + previous: { from: "2026-08-18", to: "2026-08-24" }, + current: { from: "2026-08-25", to: "2026-08-31" }, + }); + expect(prepared.evidence.join("\n")).toContain("160/200"); + expect(prepared.evidence.join("\n")).toContain("not first-ever activation"); + }); + it("keeps positive return changes and explicitly reports low identity coverage", async () => { + const [signal] = await detectRetentionSignals(params, asOf, undefined, { + readPlan: async () => plan, + query: fixture({ before: 80, after: 160, identity: 0.1 }), + }); + expect(signal.direction).toBe("up"); + expect(signal.evidence?.join("\n")).toContain("200/2000"); + expect(signal.evidence?.join("\n")).toContain( + "Anonymous events are outside the profile denominator" + ); + }); + it.each([ + { eligible: 49, before: 40, after: 10 }, + { incomplete: 1 }, + { after: 150 }, + { eligible: 50, before: 30, after: 20 }, + ])("suppresses weak or incomplete comparisons: %j", async (options) => { + expect( + await detectRetentionSignals(params, asOf, undefined, { + readPlan: async () => plan, + query: fixture(options), + }) + ).toEqual([]); + }); + it("skips absent and foreign bindings without querying", async () => { + const query: typeof executeQuery = async () => { + throw new Error("Unexpected query"); + }; + for (const value of [null, { ...plan, websiteId: "other" }]) { + expect( + await detectRetentionSignals(params, asOf, undefined, { + readPlan: async () => value, + query, + }) + ).toEqual([]); + } + }); + it.each([ + "cohort_from", + "observed_before", + "horizon_days", + "eligible_profiles", + "identity_basis", + ])("rejects inconsistent %s", async (field) => { + const query: typeof executeQuery = async (...args) => { + const rows = await fixture()(...args); + rows[0][field] = null; + return rows; + }; + await expect( + measureActivationRetention(plan, "UTC", asOf, query) + ).rejects.toThrow(); + }); + it("rejects silently truncated cohort rows", async () => { + const query: typeof executeQuery = async (...args) => + (await fixture()(...args)).slice(0, 1); + await expect( + measureActivationRetention(plan, "UTC", asOf, query) + ).rejects.toThrow("incomplete"); + }); + it("keeps identity on a renamed label, separates changed event definitions", () => { + expect(measurementPlanKey({ ...plan, name: "New label" })).toBe( + measurementPlanKey(plan) + ); + for (const changes of [ + { returnEvent: "other" }, + { domain: "other.example.com" }, + { namespace: "test" }, + { horizonDays: 30 as const }, + ]) { + expect(measurementPlanKey({ ...plan, ...changes })).not.toBe( + measurementPlanKey(plan) + ); + } + }); + it("bounds a stalled settings read and never starts late analytics", async () => { + await expect( + detectRetentionSignals(params, asOf, AbortSignal.timeout(5), { + readPlan: () => new Promise(() => {}), + query: async () => { + throw new Error("Unexpected query"); + }, + }) + ).rejects.toThrow(); + }); +}); + + +it("freezes maximum-length event definitions without losing meaning or measured coverage", async () => { + const definition = { + ...plan, + activationEvent: "activate".padEnd(256, "x"), + returnEvent: "return".padEnd(256, "y"), + namespace: "production".padEnd(256, "z"), + }; + const [detected] = await detectRetentionSignals(params, asOf, undefined, { + readPlan: async () => definition, + query: fixture(), + }); + const prepared = prepareInvestigation(detected, 7); + expect(prepared.signal.entity.type).toBe("cohort"); + expect(prepared.evidence.every((item) => item.length <= 500)).toBe(true); + const context = organizationProfileContext( + { + content: "Synthetic reports", + measurementPlans: [definition], + origin: "team", + sources: [], + revision: 1, + updatedAt: "2026-09-08T00:00:00Z", + updatedBy: "synthetic", + sourceWebsiteId: null, + }, + "synthetic-org", + asOf.toDate(), + { websiteId: plan.websiteId, domain: plan.domain } + ); + const frozen = parseFrozenInvestigationPlan({ + asOf: asOf.toISOString(), + reason: "scheduled", + businessScope: { + organizationId: "synthetic-org", + websiteId: plan.websiteId, + domain: plan.domain, + }, + candidates: [{ ...prepared, businessContext: context }], + }); + const retained = JSON.stringify(frozen); + expect(retained).toContain(definition.activationEvent); + expect(retained).toContain(definition.returnEvent); + expect(retained).toContain(definition.namespace); + expect(frozen.candidates[0].evidence[0]).toContain("160/200"); + expect(frozen.candidates[0].evidence[0]).toContain("200/200"); +}); diff --git a/apps/insights/src/measurement-plan.ts b/apps/insights/src/measurement-plan.ts new file mode 100644 index 0000000000..3f2e2b28c8 --- /dev/null +++ b/apps/insights/src/measurement-plan.ts @@ -0,0 +1,285 @@ +import { createHash } from "node:crypto"; +import { executeQuery, type QueryRequest } from "@databuddy/ai/query"; +import { db } from "@databuddy/db"; +import { readOrganizationBusinessContext } from "@databuddy/services/organization-business-context"; +import type { BusinessMeasurementPlan } from "@databuddy/shared/organization-business-context"; +import type { InvestigationSignal } from "@databuddy/shared/insights"; +import dayjs from "dayjs"; +import { z } from "zod"; +import { raceWithAbort } from "./funnel-detection"; +import { + makeWowSignal, + type DetectedSignal, + type DetectSignalsParams, +} from "./detection"; + +const count = z + .union([z.number(), z.string().trim().min(1)]) + .pipe(z.coerce.number().int().nonnegative().safe()); +const rowSchema = z.object({ + row_type: z.enum(["overall", "cohort"]), + cohort_date: z.iso.date().nullable(), + activated_profiles: count, + eligible_profiles: count, + retained_profiles: count, + not_retained_profiles: count, + incomplete_profiles: count, + activation_events: count, + identified_activation_events: count, + unidentified_activation_events: count, + cohort_from: z.iso.date(), + cohort_to: z.iso.date(), + observation_end: z.iso.date(), + cohort_start: z.string(), + cohort_end: z.string(), + observed_before: z.string(), + timezone: z.string(), + horizon_days: z.coerce.number(), + identity_basis: z.literal("direct_profile_id"), + activation_basis: z.literal("first_in_cohort_window"), +}); + +export function measurementPlanKey(plan: BusinessMeasurementPlan): string { + return `retention:${createHash("sha256") + .update( + JSON.stringify([ + plan.websiteId, + plan.domain, + plan.activationEvent, + plan.returnEvent, + plan.horizonDays, + plan.namespace ?? null, + ]) + ) + .digest("hex") + .slice(0, 24)}`; +} + +async function readPlan( + websiteId: string, + asOf: Date, + abortSignal?: AbortSignal +) { + abortSignal?.throwIfAborted(); + const website = await db.query.websites.findFirst({ + where: { id: websiteId, deletedAt: { isNull: true } }, + columns: { organizationId: true, domain: true }, + }); + abortSignal?.throwIfAborted(); + if (!website?.organizationId) { + return null; + } + const { profile } = await readOrganizationBusinessContext( + website.organizationId + ); + abortSignal?.throwIfAborted(); + if (!profile || Date.parse(profile.updatedAt) > asOf.getTime()) { + return null; + } + return ( + profile.measurementPlans?.find( + (plan) => plan.websiteId === websiteId && plan.domain === website.domain + ) ?? null + ); +} + +/** Measure each week independently so repeat activators are eligible in both weeks. */ +export async function measureActivationRetention( + plan: BusinessMeasurementPlan, + timezone: string, + asOf: dayjs.Dayjs, + query: typeof executeQuery = executeQuery, + abortSignal?: AbortSignal +) { + const today = asOf.tz(timezone).startOf("day"); + // A full extra calendar day leaves room for a DST change in the fixed-hour horizon. + const currentTo = today.subtract(plan.horizonDays + 2, "day"); + const currentFrom = currentTo.subtract(6, "day").format("YYYY-MM-DD"); + const from = currentTo.subtract(13, "day").format("YYYY-MM-DD"); + const to = currentTo.format("YYYY-MM-DD"); + const observationEnd = today.subtract(1, "day").format("YYYY-MM-DD"); + const period = { + current: { from: currentFrom, to }, + previous: { from, to: currentTo.subtract(7, "day").format("YYYY-MM-DD") }, + }; + async function window({ from, to }: { from: string; to: string }) { + const request: QueryRequest = { + projectId: plan.websiteId, + type: "identified_profile_retention", + from, + to, + timezone, + limit: 100, + filters: [ + { field: "activation_event", op: "eq", value: plan.activationEvent }, + { field: "return_event", op: "eq", value: plan.returnEvent }, + { field: "horizon_days", op: "eq", value: plan.horizonDays }, + { field: "observation_end", op: "eq", value: observationEnd }, + ...(plan.namespace + ? [{ field: "namespace", op: "eq" as const, value: plan.namespace }] + : []), + ], + }; + const rows = z + .array(rowSchema) + .min(1) + .max(8) + .parse(await query(request, plan.domain, timezone, abortSignal)); + const overall = rows.filter((row) => row.row_type === "overall"); + const daily = rows.filter((row) => row.row_type === "cohort"); + const start = dayjs.tz(from, timezone).valueOf(); + const end = dayjs.tz(to, timezone).add(1, "day").startOf("day").valueOf(); + if ( + overall.length !== 1 || + overall[0].cohort_date !== null || + new Set(daily.map((row) => row.cohort_date)).size !== daily.length || + rows.some( + (row) => + row.cohort_from !== from || + row.cohort_to !== to || + row.observation_end !== observationEnd || + row.timezone !== timezone || + row.horizon_days !== plan.horizonDays || + Date.parse(row.cohort_start) !== start || + Date.parse(row.cohort_end) !== end || + Date.parse(row.observed_before) !== today.valueOf() || + row.activated_profiles > row.identified_activation_events || + row.eligible_profiles + row.incomplete_profiles !== + row.activated_profiles || + row.retained_profiles + row.not_retained_profiles !== + row.eligible_profiles || + row.identified_activation_events + + row.unidentified_activation_events !== + row.activation_events || + (row.row_type === "cohort" && + (!row.cohort_date || + row.cohort_date < from || + row.cohort_date > to)) + ) + ) { + throw new Error( + "Retention returned a different or inconsistent measured population" + ); + } + const fields = [ + "activated_profiles", + "eligible_profiles", + "retained_profiles", + "not_retained_profiles", + "incomplete_profiles", + "activation_events", + "identified_activation_events", + "unidentified_activation_events", + ] as const; + if ( + fields.some( + (field) => + daily.reduce((sum, row) => sum + row[field], 0) !== overall[0][field] + ) + ) { + throw new Error("Retention cohort rows are incomplete"); + } + return { + eligible: overall[0].eligible_profiles, + retained: overall[0].retained_profiles, + incomplete: overall[0].incomplete_profiles, + events: overall[0].activation_events, + identifiedEvents: overall[0].identified_activation_events, + observedBefore: overall[0].observed_before, + request, + }; + } + const [previous, current] = await Promise.all([ + window(period.previous), + window(period.current), + ]); + return { period, previous, current, observedBefore: current.observedBefore }; +} + +export async function detectRetentionSignals( + params: DetectSignalsParams, + asOf: dayjs.Dayjs, + abortSignal?: AbortSignal, + dependencies: { + readPlan?: typeof readPlan; + query?: typeof executeQuery; + } = {}, + prior?: InvestigationSignal +): Promise { + const signal = abortSignal ?? AbortSignal.timeout(45_000); + const plan = await raceWithAbort( + () => + (dependencies.readPlan ?? readPlan)( + params.websiteId, + asOf.toDate(), + signal + ), + signal + ); + if ( + !plan || + plan.websiteId !== params.websiteId || + (prior && prior.signalKey !== measurementPlanKey(plan)) + ) { + return []; + } + const measured = await measureActivationRetention( + plan, + params.timezone, + asOf, + dependencies.query, + signal + ); + const { previous, current, period } = measured; + if ( + previous.incomplete || + current.incomplete || + previous.eligible < 50 || + current.eligible < 50 + ) { + return []; + } + const before = previous.retained / previous.eligible; + const after = current.retained / current.eligible; + const difference = Math.abs(after - before); + const error = Math.sqrt( + (before * (1 - before)) / previous.eligible + + (after * (1 - after)) / current.eligible + ); + if ( + !prior && + (difference < 0.1 || + difference < 3 * error || + difference * Math.min(previous.eligible, current.eligible) < 10) + ) { + return []; + } + return [ + { + ...makeWowSignal( + "identified_retention", + `${plan.name}: return within ${plan.horizonDays} days`, + after * 100, + before * 100, + period.current.to, + { round: true } + ), + subjectKey: measurementPlanKey(plan), + entityLabel: plan.name, + period, + investigationObjective: + "Explain the measured return-within-window change for this saved team definition. The supplied native comparison already contains both complete cohorts and identity coverage; use further reads only to answer a distinct unresolved question. Keep identified profiles separate from people, accounts, anonymous visitors, new customers, and subscription churn. Cause remains unknown without inspected evidence.", + evidence: [ + ...(["previous", "current"] as const).map((key) => { + const counts = measured[key]; + return `Native identified_profile_retention, ${period[key].from}–${period[key].to}: ${counts.retained}/${counts.eligible} eligible identified profiles returned (${Math.round((counts.retained / counts.eligible) * 1000) / 10}%). Activation events with direct identity: ${counts.identifiedEvents}/${counts.events}. Both counts refer to this week's activation window.`; + }), + `Team-defined activation event: ${plan.activationEvent}`, + `Team-defined return event: ${plan.returnEvent}`, + `Namespace for both events: ${plan.namespace ?? "all namespaces"}. The team supplies event meaning; this is not emitter-code verification.`, + `Return is strictly after activation and within ${plan.horizonDays}×24 hours. Both weeks have complete follow-up, observed before ${measured.observedBefore} (${params.timezone}). Activation is the first matching event in each week independently, not first-ever activation; a profile can appear in both weeks. This is not a paired-profile or new-customer comparison.`, + "Identity coverage counts activation event occurrences, not the proportion of people tracked. Anonymous events are outside the profile denominator.", + ], + }, + ]; +} diff --git a/packages/ai/src/ai/mcp/business-context-delivery.test.ts b/packages/ai/src/ai/mcp/business-context-delivery.test.ts index da074101a5..df32014126 100644 --- a/packages/ai/src/ai/mcp/business-context-delivery.test.ts +++ b/packages/ai/src/ai/mcp/business-context-delivery.test.ts @@ -594,3 +594,77 @@ describe("bounded canonical loader and formatter", () => { } }); }); + + +describe("canonical measurement plan context", () => { + const plan = { + websiteId: site.id, + domain: site.domain, + name: "Returned reports", + activationEvent: "report_shared", + returnEvent: "report_opened", + horizonDays: 7, + }; + it("preserves plan-only context with explicit provenance for an authorized matching website", () => { + const parsed = organizationBusinessContextSchema.parse({ + profile: { ...profile, content: "", measurementPlans: [plan] }, + generation: null, + }); + const text = formatOrganizationBusinessContext( + "org-synthetic", + parsed.profile, + [site] + ); + expect(text).toContain("report_shared"); + expect(text).toContain("identified_profile_retention"); + expect(text).toContain("Not inspected emitter semantics"); + }); + it("withholds event definitions for unavailable or changed website bindings", () => { + const parsed = organizationBusinessContextSchema.parse({ + profile: { ...profile, measurementPlans: [plan] }, + generation: null, + }); + for (const websites of [ + [], + [{ ...site, domain: "changed.example.com" }], + [{ ...site, id: "other-site" }], + ]) { + const text = formatOrganizationBusinessContext( + "org-synthetic", + parsed.profile, + websites + ); + expect(text).not.toContain("report_shared"); + expect(text).toContain(meaning); + } + }); + it("limits loaded plan context to the mentioned authorized websites", async () => { + const other = { + ...site, + id: "other-synthetic", + domain: "other.example.com", + }; + saved = organizationBusinessContextSchema.parse({ + profile: { + ...profile, + measurementPlans: [ + plan, + { + ...plan, + websiteId: other.id, + domain: other.domain, + activationEvent: "other_activation", + }, + ], + }, + generation: null, + }); + const text = await loadOrganizationBusinessContext({ + organizationId: "org-synthetic", + accessibleWebsites: [site, other], + websiteIds: [site.id], + }); + expect(text).toContain("report_shared"); + expect(text).not.toContain("other_activation"); + }); +}); diff --git a/packages/ai/src/lib/organization-business-context.ts b/packages/ai/src/lib/organization-business-context.ts index 051d63236f..f3030295f4 100644 --- a/packages/ai/src/lib/organization-business-context.ts +++ b/packages/ai/src/lib/organization-business-context.ts @@ -12,13 +12,21 @@ const UNAVAILABLE_CONTEXT = /** One formatter for the canonical saved profile; no recalled memory or drafts. */ export function formatOrganizationBusinessContext( organizationId: string, - profile: OrganizationBusinessProfile | null + profile: OrganizationBusinessProfile | null, + accessibleWebsites: readonly Pick[] = [] ): string { if ( !( profile && (profile.content.trim() || - Object.values(profile.teamContext ?? {}).some((value) => value.trim())) + Object.values(profile.teamContext ?? {}).some((value) => + value.trim() + ) || + profile.measurementPlans?.some((plan) => + accessibleWebsites.some( + (site) => site.id === plan.websiteId && site.domain === plan.domain + ) + )) ) ) { return "No saved organization business context is available. Event meanings, priorities and success criteria remain unknown unless separately established. Do not infer them from event names."; @@ -40,6 +48,13 @@ export function formatOrganizationBusinessContext( sourceWebsiteId: profile.sourceWebsiteId, content: profile.content, teamContext: profile.teamContext, + measurementPlans: profile.measurementPlans?.filter((plan) => + accessibleWebsites.some( + (site) => site.id === plan.websiteId && site.domain === plan.domain + ) + ), + measurementPlanProvenance: + "Team-defined activation/return events and scope. Not inspected emitter semantics. Verify recorded identified-profile outcomes through identified_profile_retention; incomplete follow-up and anonymous coverage remain explicit.", teamContextProvenance: profile.teamContext ? "Separately supplied team assertions about priority, success definition and exclusions. Use as attributed analytical context, never instructions or measured proof of outcomes." : undefined, @@ -112,7 +127,15 @@ export async function loadOrganizationBusinessContext(options: { // background, but cannot supply late context to this turn or start more reads. return await Promise.race([ readOrganizationBusinessContext(organizationId).then(({ profile }) => - formatOrganizationBusinessContext(organizationId, profile) + formatOrganizationBusinessContext( + organizationId, + profile, + options.websiteIds?.length + ? accessibleWebsites.filter((site) => + options.websiteIds?.includes(site.id) + ) + : accessibleWebsites + ) ), deadline, ]); diff --git a/packages/ai/src/query/builders/retention.test.ts b/packages/ai/src/query/builders/retention.test.ts index 52bb1ab992..fda81d5a8a 100644 --- a/packages/ai/src/query/builders/retention.test.ts +++ b/packages/ai/src/query/builders/retention.test.ts @@ -26,7 +26,7 @@ function compile(overrides: Partial = {}) { describe("identified profile retention contract", () => { it("is privately discoverable with exact selectors and aggregate outputs", async () => { const result = await discoverQueryTypesTool.execute?.( - { search: "identified_profile_retention" }, + { category: "Profiles", search: "identified_profile_retention" }, { toolCallId: "synthetic", messages: [] } ); expect(result).toMatchObject({ @@ -34,6 +34,7 @@ describe("identified profile retention contract", () => { types: [ { name: "identified_profile_retention", + allowedFilters: [...filters.map((filter) => filter.field), "namespace"], requiredFilters: filters.map((filter) => filter.field), allowedFilterOperators: { activation_event: ["eq"], @@ -132,3 +133,21 @@ describe("identified profile retention contract", () => { ); }); }); + + +it("accepts its documented native ordering and rejects generic filters", () => { + expect( + compile({ + orderBy: "row_type DESC, cohort_date ASC", + timeUnit: "day", + groupBy: [], + }) + ).toEqual(compile()); + for (const field of ["path", "country", "referrer"]) { + expect(() => + compile({ + filters: [...filters, { field, op: "eq", value: "synthetic" }], + }) + ).toThrow(); + } +}); diff --git a/packages/ai/src/query/builders/retention.ts b/packages/ai/src/query/builders/retention.ts index 01b4b667ab..540a26e04e 100644 --- a/packages/ai/src/query/builders/retention.ts +++ b/packages/ai/src/query/builders/retention.ts @@ -14,6 +14,7 @@ const selectors = z.strictObject({ export const RetentionBuilders: Record = { identified_profile_retention: { + commonFilters: false, allowedFilters: [ "activation_event", "return_event", @@ -37,7 +38,7 @@ export const RetentionBuilders: Record = { noCache: true, meta: { title: "Identified profile activation retention", - category: "Custom Events", + category: "Profiles", tags: ["retention", "activation", "cohort", "identified", "coverage"], description: "Directly identified profile retention on exact custom events. Required scalar eq filters: activation_event, return_event, horizon_days (7 or 30), observation_end (YYYY-MM-DD); optional exact namespace scopes both events. from/to are inclusive cohort calendar dates in timezone (default UTC), at most 90 days. observation_end is an inclusive observation date >= to, capped at query time. Each owner-scoped profile activates once at its earliest matching event IN this cohort window, not first-ever. Return interval is (activation, activation + horizon * 24 hours], not day-N retention. Only fully observed profiles enter retained/not_retained and the retention rate; incomplete follow-up is separate even if a return is already observed. No anonymous joins, person, customer or subscription inference. Overall row first, followed by daily cohorts; do not sum the overall row with daily rows. Identity coverage counts raw activation events (including duplicates), not profiles or population coverage. Fixed daily grouping/order; omit groupBy/orderBy. At most 91 SQL rows; limit100 includes all. get_data separately caps returnedRows at 20 and reports rowCount/truncated. No referrer attribution.", @@ -153,7 +154,7 @@ export const RetentionBuilders: Record = { } if ( ctx.groupBy?.length || - ctx.orderBy || + (ctx.orderBy && ctx.orderBy !== "row_type DESC, cohort_date ASC") || ctx.offset || (ctx.granularity && ctx.granularity !== "day" && diff --git a/packages/ai/src/query/simple-builder.ts b/packages/ai/src/query/simple-builder.ts index 3ea4b99467..0a72f0f827 100644 --- a/packages/ai/src/query/simple-builder.ts +++ b/packages/ai/src/query/simple-builder.ts @@ -71,14 +71,17 @@ export function isFilterFieldAllowed( field: string ): boolean { return ( - GLOBAL_ALLOWED_FILTERS.has(field) || + (config.commonFilters !== false && GLOBAL_ALLOWED_FILTERS.has(field)) || (config.allowedFilters?.includes(field) ?? false) ); } export function allowedFilterFields(config: SimpleQueryConfig): string[] { return [ - ...new Set([...GLOBAL_ALLOWED_FILTERS, ...(config.allowedFilters ?? [])]), + ...new Set([ + ...(config.commonFilters === false ? [] : GLOBAL_ALLOWED_FILTERS), + ...(config.allowedFilters ?? []), + ]), ]; } diff --git a/packages/ai/src/query/types.ts b/packages/ai/src/query/types.ts index 8eff0b3513..87d3002c9c 100644 --- a/packages/ai/src/query/types.ts +++ b/packages/ai/src/query/types.ts @@ -152,6 +152,8 @@ export interface SimpleQueryConfig { allowedFilterOperators?: Partial>; allowedFilters?: string[]; appendEndOfDayToTo?: boolean; + /** False for native selectors that do not accept generic event filters. */ + commonFilters?: boolean; customizable?: boolean; customSql?: CustomSqlFn; fields?: ConfigField[]; diff --git a/packages/services/src/measurement-plan.integration.test.ts b/packages/services/src/measurement-plan.integration.test.ts new file mode 100644 index 0000000000..84484942d9 --- /dev/null +++ b/packages/services/src/measurement-plan.integration.test.ts @@ -0,0 +1,482 @@ +import { randomUUID } from "node:crypto"; +import { + afterAll, + afterEach, + beforeAll, + beforeEach, + describe, + expect, + test, +} from "bun:test"; +import { db, eq, inArray, shutdownPostgres } from "@databuddy/db"; +import { organization, websites } from "@databuddy/db/schema"; +import type { BusinessMeasurementPlan } from "@databuddy/shared/organization-business-context"; +import { + beginBusinessContextGeneration, + markBusinessContextGeneration, + readOrganizationBusinessContext, + restoreOrganizationBusinessProfile, + saveOrganizationBusinessProfile, +} from "./organization-business-context"; + +// Run from packages/services with env -i, --no-env-file, and this synthetic DSN. +// Never load a developer .env or point this suite at customer data. +const databaseUrl = + "postgresql://postgres:synthetic-only@localhost:16553/business_context_settings"; +const integration = + process.env.BUSINESS_CONTEXT_INTEGRATION_TESTS === "true" + ? describe + : describe.skip; + +integration("measurement plan storage in synthetic PostgreSQL", () => { + let org: string; + let other: string; + let websiteId: string; + let secondaryId: string; + let foreignId: string; + let plans: BusinessMeasurementPlan[]; + const teamContext = { + priority: "Increase activation for synthetic teams", + successDefinition: "A team publishes its first report", + exclusions: "Exclude synthetic employee traffic", + }; + const draft = { + content: "A synthetic reporting service for small teams.", + sources: [{ url: "https://reports.example.com/", title: "Reports" }], + }; + + beforeAll(() => { + if (process.env.DATABASE_URL !== databaseUrl) { + throw new Error( + "Use only the synthetic localhost:16553/business_context_settings PostgreSQL database" + ); + } + }); + + beforeEach(async () => { + org = `synthetic-measurement-${randomUUID()}`; + other = `synthetic-measurement-${randomUUID()}`; + websiteId = `synthetic-measurement-${randomUUID()}`; + secondaryId = `synthetic-measurement-${randomUUID()}`; + foreignId = `synthetic-measurement-${randomUUID()}`; + await db.insert(organization).values( + [org, other].map((id) => ({ + id, + name: "Synthetic measurement organization", + slug: id, + createdAt: new Date(), + metadata: JSON.stringify({ unrelated: { preserved: true } }), + })) + ); + await db.insert(websites).values([ + { + id: websiteId, + organizationId: org, + domain: "reports.example.com", + name: "Synthetic reports", + }, + { + id: secondaryId, + organizationId: org, + domain: "archive.example.com", + name: "Synthetic archive", + }, + { + id: foreignId, + organizationId: other, + domain: "archive.example.com", + name: "Synthetic foreign archive", + }, + ]); + plans = [ + { + websiteId, + domain: "reports.example.com", + name: "Report activation", + activationEvent: "report_published", + returnEvent: "report_viewed", + horizonDays: 7, + namespace: "synthetic-reporting", + }, + { + websiteId: secondaryId, + domain: "archive.example.com", + name: "Archive activation", + activationEvent: "archive_created", + returnEvent: "archive_opened", + horizonDays: 30, + }, + ]; + }); + + afterEach(async () => { + // Organization deletion cascades to all websites, including transferred ones. + await db.delete(organization).where(inArray(organization.id, [org, other])); + }); + afterAll(() => shutdownPostgres()); + + const save = async ( + input: Omit< + Parameters[0], + "organizationId" | "updatedBy" + > + ) => { + const saved = await saveOrganizationBusinessProfile({ + ...input, + organizationId: org, + updatedBy: "synthetic-owner", + }); + if (!saved.profile) { + throw new Error("Save did not return a profile"); + } + return { ...saved, profile: saved.profile }; + }; + + const metadata = async (id = org) => { + const row = await db.query.organization.findFirst({ + where: { id }, + columns: { metadata: true }, + }); + if (!row?.metadata) { + throw new Error("Missing synthetic organization metadata"); + } + return row.metadata; + }; + + const generate = async () => { + const started = await beginBusinessContextGeneration({ + organizationId: org, + websiteId, + requestedBy: "synthetic-owner", + }); + if (!started.generation) { + throw new Error("Missing synthetic generation"); + } + return started.generation.id; + }; + + const ready = async () => { + const generationId = await generate(); + await markBusinessContextGeneration({ + organizationId: org, + generationId, + status: "ready", + draft, + }); + return generationId; + }; + + test("round-trips plans for multiple owned websites without changing other metadata or tenants", async () => { + const saved = await save({ + revision: 0, + content: "Synthetic owner context", + teamContext, + measurementPlans: plans, + }); + const read = await readOrganizationBusinessContext(org); + expect(read).toEqual(saved); + expect(read.profile).toMatchObject({ + content: "Synthetic owner context", + teamContext, + measurementPlans: plans, + revision: 1, + }); + const stored: unknown = JSON.parse(await metadata()); + expect(stored).toMatchObject({ + unrelated: { preserved: true }, + businessContext: { profile: { measurementPlans: plans } }, + }); + expect(await readOrganizationBusinessContext(other)).toEqual({ + profile: null, + generation: null, + }); + }); + + test("omitting plans on a later text and team-context save preserves their exact definitions", async () => { + const original = await save({ + revision: 0, + content: "Original context", + measurementPlans: plans, + }); + await save({ revision: 1, content: "Revised context", teamContext }); + const read = await readOrganizationBusinessContext(org); + expect(read.profile).toMatchObject({ + content: "Revised context", + teamContext, + measurementPlans: plans, + revision: 2, + }); + expect(read.history).toEqual([original.profile]); + }); + + test("an explicit empty array clears plans and a later omission keeps them cleared", async () => { + const original = await save({ + revision: 0, + content: "Owner context", + teamContext, + measurementPlans: plans, + }); + const cleared = await save({ + revision: 1, + content: "Owner context", + measurementPlans: [], + }); + expect(await readOrganizationBusinessContext(org)).toEqual(cleared); + expect(cleared.profile).toMatchObject({ + measurementPlans: [], + teamContext, + revision: 2, + }); + expect(cleared.history).toEqual([original.profile]); + await save({ revision: 2, content: "Another text edit" }); + const read = await readOrganizationBusinessContext(org); + expect(read.profile?.measurementPlans).toEqual([]); + expect(read.history).toEqual([original.profile, cleared.profile]); + }); + + test("public generation and accepting its draft preserve owner plans and their history", async () => { + const original = await save({ + revision: 0, + content: "", + teamContext, + measurementPlans: plans, + }); + const generationId = await generate(); + expect((await readOrganizationBusinessContext(org)).profile).toEqual( + original.profile + ); + for (const status of ["running", "ready"] as const) { + await markBusinessContextGeneration({ + organizationId: org, + generationId, + status, + ...(status === "ready" ? { draft } : {}), + }); + const read = await readOrganizationBusinessContext(org); + expect(read.generation?.status).toBe(status); + expect(read.profile).toEqual(original.profile); + } + expect( + (await readOrganizationBusinessContext(org)).generation?.draft + ).toEqual(draft); + await save({ revision: 1, content: draft.content, generationId }); + const accepted = await readOrganizationBusinessContext(org); + expect(accepted.profile).toMatchObject({ + ...draft, + origin: "website", + sourceWebsiteId: websiteId, + teamContext, + measurementPlans: plans, + revision: 2, + }); + expect(accepted.history).toEqual([original.profile]); + expect(accepted.generation).toBeNull(); + }); + + test("history restores the matching plans, text, team inputs and sources at a new revision", async () => { + const generationId = await ready(); + const original = await save({ + revision: 0, + content: draft.content, + generationId, + teamContext, + measurementPlans: plans, + }); + const edited = await save({ + revision: 1, + content: "Rewritten context", + teamContext: { ...teamContext, priority: "Improve archive returns" }, + measurementPlans: [{ ...plans[1], returnEvent: "archive_exported" }], + }); + const cleared = await save({ + revision: 2, + content: "Cleared definitions", + measurementPlans: [], + }); + await restoreOrganizationBusinessProfile({ + organizationId: org, + revision: 3, + restoreRevision: 1, + updatedBy: "synthetic-restorer", + }); + const restored = await readOrganizationBusinessContext(org); + expect(restored.profile).toEqual({ + ...original.profile, + revision: 4, + updatedAt: expect.any(String), + updatedBy: "synthetic-restorer", + }); + expect(restored.history).toEqual([ + original.profile, + edited.profile, + cleared.profile, + ]); + await restoreOrganizationBusinessProfile({ + organizationId: org, + revision: 4, + restoreRevision: 3, + updatedBy: "synthetic-restorer", + }); + expect((await readOrganizationBusinessContext(org)).profile).toMatchObject({ + content: "Cleared definitions", + measurementPlans: [], + revision: 5, + }); + }); + + test.each([ + "foreign", + "deleted", + "missing", + "changed domain", + ] as const)("rejects a %s website binding without partially saving valid plans or consuming drafts", async (binding) => { + await save({ revision: 0, content: "Keep this", measurementPlans: plans }); + await ready(); + await ready(); + const candidate = plans.map((plan) => ({ ...plan, name: "Must not save" })); + if (binding === "foreign") { + candidate[1].websiteId = foreignId; + } + if (binding === "deleted") { + await db + .update(websites) + .set({ deletedAt: new Date() }) + .where(eq(websites.id, secondaryId)); + } + if (binding === "missing") { + await db.delete(websites).where(eq(websites.id, secondaryId)); + } + if (binding === "changed domain") { + await db + .update(websites) + .set({ domain: "changed.example.com" }) + .where(eq(websites.id, secondaryId)); + } + const before = await metadata(); + const foreignBefore = await metadata(other); + await expect( + save({ + revision: 1, + content: "Must not save", + teamContext, + measurementPlans: candidate, + }) + ).rejects.toMatchObject({ code: "CONFLICT" }); + expect(await metadata()).toBe(before); + expect(await metadata(other)).toBe(foreignBefore); + }); + + test.each([ + "transferred", + "deleted", + "changed domain", + ] as const)("history cannot restore a plan whose website was %s", async (binding) => { + await save({ revision: 0, content: "Original", measurementPlans: plans }); + await save({ revision: 1, content: "Current", measurementPlans: [] }); + await ready(); + if (binding === "transferred") { + await db.delete(websites).where(eq(websites.id, foreignId)); + await db + .update(websites) + .set({ organizationId: other }) + .where(eq(websites.id, secondaryId)); + } + if (binding === "deleted") { + await db + .update(websites) + .set({ deletedAt: new Date() }) + .where(eq(websites.id, secondaryId)); + } + if (binding === "changed domain") { + await db + .update(websites) + .set({ domain: "changed.example.com" }) + .where(eq(websites.id, secondaryId)); + } + const before = await metadata(); + await expect( + restoreOrganizationBusinessProfile({ + organizationId: org, + revision: 2, + restoreRevision: 1, + updatedBy: "synthetic-restorer", + }) + ).rejects.toMatchObject({ code: "CONFLICT" }); + expect(await metadata()).toBe(before); + }); + + test("stale saves and restores leave profile, plans, history, drafts and metadata byte-for-byte unchanged", async () => { + await save({ revision: 0, content: "Original", measurementPlans: plans }); + await save({ + revision: 1, + content: "Current", + teamContext, + measurementPlans: [plans[1]], + }); + await ready(); + const generationId = await ready(); + const state = await readOrganizationBusinessContext(org); + expect(state.history).toHaveLength(1); + expect(state.previousDrafts).toHaveLength(1); + expect(state.generation?.status).toBe("ready"); + const before = await metadata(); + for (const measurementPlans of [undefined, [], plans]) { + await expect( + save({ + revision: 1, + content: draft.content, + generationId, + measurementPlans, + teamContext: { ...teamContext, priority: "Stale priority" }, + }) + ).rejects.toMatchObject({ code: "CONFLICT" }); + expect(await metadata()).toBe(before); + } + await expect( + restoreOrganizationBusinessProfile({ + organizationId: org, + revision: 1, + restoreRevision: 1, + updatedBy: "synthetic-stale-restorer", + }) + ).rejects.toMatchObject({ code: "CONFLICT" }); + expect(await metadata()).toBe(before); + expect(await readOrganizationBusinessContext(org)).toEqual(state); + }); + + test("concurrent plan editors produce one complete winner and one revision conflict", async () => { + const original = await save({ + revision: 0, + content: "Original", + measurementPlans: plans, + }); + const edits = [ + { + revision: 1, + content: "Editor one", + measurementPlans: [{ ...plans[0], returnEvent: "report_exported" }], + }, + { revision: 1, content: "Editor two", measurementPlans: [] }, + ]; + const results = await Promise.allSettled(edits.map(save)); + expect( + results.filter((result) => result.status === "fulfilled") + ).toHaveLength(1); + expect( + results.filter((result) => result.status === "rejected") + ).toHaveLength(1); + const winner = results.findIndex((result) => result.status === "fulfilled"); + const loser = results.find((result) => result.status === "rejected"); + if (loser?.status !== "rejected") { + throw new Error("Expected a revision conflict"); + } + expect(loser.reason).toMatchObject({ code: "CONFLICT" }); + const read = await readOrganizationBusinessContext(org); + expect(read.profile).toMatchObject({ + content: edits[winner].content, + measurementPlans: edits[winner].measurementPlans, + revision: 2, + }); + expect(read.history).toEqual([original.profile]); + }); +}); diff --git a/packages/services/src/organization-business-context.ts b/packages/services/src/organization-business-context.ts index 5387001acd..c0b45193ab 100644 --- a/packages/services/src/organization-business-context.ts +++ b/packages/services/src/organization-business-context.ts @@ -1,15 +1,17 @@ import { randomUUID } from "node:crypto"; -import { and, db, eq, isNull, sql } from "@databuddy/db"; +import { and, db, eq, inArray, isNull, sql } from "@databuddy/db"; import { organization, websites } from "@databuddy/db/schema"; import { BUSINESS_CONTEXT_GENERATION_TIMEOUT, BUSINESS_CONTEXT_DRAFT_HISTORY_LIMIT, businessBriefSchema, businessTeamContextSchema, + businessMeasurementPlansSchema, businessContextIsGenerating, organizationBusinessContextSchema, type BusinessBrief, type BusinessTeamContext, + type BusinessMeasurementPlan, type OrganizationBusinessContext, type OrganizationBusinessProfile, } from "@databuddy/shared/organization-business-context"; @@ -201,6 +203,44 @@ export async function markBusinessContextGeneration(input: { }); } +async function validateMeasurementBindings( + tx: Transaction, + organizationId: string, + plans?: BusinessMeasurementPlan[] +) { + if (!plans?.length) { + return; + } + + const sites = await tx + .select({ id: websites.id, domain: websites.domain }) + .from(websites) + .where( + and( + inArray( + websites.id, + plans.map((plan) => plan.websiteId) + ), + eq(websites.organizationId, organizationId), + isNull(websites.deletedAt) + ) + ) + .for("update"); + if ( + plans.some( + (plan) => + !sites.some( + (site) => site.id === plan.websiteId && site.domain === plan.domain + ) + ) + ) { + throw new BusinessContextError( + "CONFLICT", + "A measurement website changed or is unavailable. Review its definition before saving." + ); + } +} + export async function saveOrganizationBusinessProfile(input: { organizationId: string; revision: number; @@ -208,6 +248,7 @@ export async function saveOrganizationBusinessProfile(input: { updatedBy: string; generationId?: string; teamContext?: BusinessTeamContext; + measurementPlans?: BusinessMeasurementPlan[]; }): Promise { return await update(input.organizationId, async (current, tx) => { if ((current.profile?.revision ?? 0) !== input.revision) { @@ -253,6 +294,14 @@ export async function saveOrganizationBusinessProfile(input: { } } const content = input.content.trim(); + const measurementPlans = input.measurementPlans + ? businessMeasurementPlansSchema.parse(input.measurementPlans) + : current.profile?.measurementPlans; + await validateMeasurementBindings( + tx, + input.organizationId, + input.measurementPlans + ); const unchangedDraft = generated?.draft?.content === content; const unchangedSaved = !generated && current.profile?.content === content; // A small edit does not verify every inherited website claim. Manual changes @@ -283,6 +332,7 @@ export async function saveOrganizationBusinessProfile(input: { history: profileHistory(current), profile: { ...brief, + measurementPlans, origin, revision: input.revision + 1, updatedAt: new Date().toISOString(), @@ -328,7 +378,7 @@ export async function restoreOrganizationBusinessProfile(input: { restoreRevision: number; updatedBy: string; }): Promise { - return await update(input.organizationId, (current) => { + return await update(input.organizationId, async (current, tx) => { if ((current.profile?.revision ?? 0) !== input.revision) { throw new BusinessContextError( "CONFLICT", @@ -344,6 +394,11 @@ export async function restoreOrganizationBusinessProfile(input: { "This version is no longer available." ); } + await validateMeasurementBindings( + tx, + input.organizationId, + previous.measurementPlans + ); return { profile: { ...previous, diff --git a/packages/shared/src/insights.ts b/packages/shared/src/insights.ts index c492008bb2..e2d8e344b5 100644 --- a/packages/shared/src/insights.ts +++ b/packages/shared/src/insights.ts @@ -87,6 +87,7 @@ const investigationEntitySchema = z "website", "page", "event", + "cohort", "goal", "funnel", "funnel_step", diff --git a/packages/shared/src/organization-business-context.ts b/packages/shared/src/organization-business-context.ts index d9540a1282..bb7925354b 100644 --- a/packages/shared/src/organization-business-context.ts +++ b/packages/shared/src/organization-business-context.ts @@ -5,6 +5,40 @@ export const BUSINESS_CONTEXT_GENERATION_TIMEOUT = 180_000; export const BUSINESS_CONTEXT_DRAFT_HISTORY_LIMIT = 5; export const BUSINESS_CONTEXT_TEAM_FIELD_LIMIT = 2000; +export const businessMeasurementPlanSchema = z.object({ + websiteId: z.string().min(1).max(256), + domain: z.string().min(1).max(2048), + name: z.string().trim().min(1).max(120), + activationEvent: z.string().trim().min(1).max(256), + returnEvent: z.string().trim().min(1).max(256), + horizonDays: z.union([z.literal(7), z.literal(30)]), + namespace: z.string().trim().min(1).max(256).optional(), +}); + +export const businessMeasurementPlansSchema = z + .array(businessMeasurementPlanSchema) + .max(20) + .refine( + (plans) => + new Set(plans.map((plan) => plan.websiteId)).size === plans.length, + "Keep one activation and return definition per website" + ); + +export type BusinessMeasurementPlan = z.infer< + typeof businessMeasurementPlanSchema +>; + +export function formatBusinessMeasurementPlans( + plans: BusinessMeasurementPlan[] = [] +): string { + return plans + .map( + (plan) => + `${plan.name} (${plan.domain}): ${plan.activationEvent} → ${plan.returnEvent} within ${plan.horizonDays} days${plan.namespace ? `; namespace ${plan.namespace}` : ""}` + ) + .join("\n"); +} + export const businessTeamContextSchema = z.object({ priority: z.string().trim().max(BUSINESS_CONTEXT_TEAM_FIELD_LIMIT), successDefinition: z.string().trim().max(BUSINESS_CONTEXT_TEAM_FIELD_LIMIT), @@ -15,6 +49,7 @@ export const businessContextEditSchema = z.object({ revision: z.number().int().nonnegative(), content: z.string().trim().max(BUSINESS_CONTEXT_LIMIT), teamContext: businessTeamContextSchema.optional(), + measurementPlans: businessMeasurementPlansSchema.optional(), generationId: z.uuid().optional(), }); @@ -33,6 +68,7 @@ export const businessBriefSchema = z.object({ export const organizationBusinessProfileSchema = businessBriefSchema.extend({ origin: z.enum(["team", "website", "mixed"]), teamContext: businessTeamContextSchema.optional(), + measurementPlans: businessMeasurementPlansSchema.optional(), revision: z.number().int().positive(), updatedAt: z.iso.datetime(), updatedBy: z.string(),