diff --git a/packages/backend/src/adapters/de/lexware-office.json b/packages/backend/src/adapters/de/lexware-office.json index 62e92645..73f5f6cd 100644 --- a/packages/backend/src/adapters/de/lexware-office.json +++ b/packages/backend/src/adapters/de/lexware-office.json @@ -109,11 +109,11 @@ "properties": { "voucherType": { "type": "string", - "description": "Comma-separated types: salesinvoice, salescreditnote, purchaseinvoice, purchasecreditnote, invoice, creditnote, orderconfirmation, quotation, deliverynote, downpaymentinvoice." + "description": "Required by Lexware. Comma-separated types: salesinvoice, salescreditnote, purchaseinvoice, purchasecreditnote, invoice, creditnote, orderconfirmation, quotation, deliverynote, downpaymentinvoice, or any." }, "voucherStatus": { "type": "string", - "description": "Comma-separated statuses: draft, open, paid, paidoff, voided, transferred, sepadebit, overdue, accepted, rejected." + "description": "Required by Lexware. Comma-separated statuses: draft, open, paid, paidoff, voided, transferred, sepadebit, overdue, accepted, rejected, or any. overdue cannot be combined with other statuses: ask for it on its own (Lexware answers 400 otherwise)." }, "archived": { "type": "boolean", @@ -139,7 +139,11 @@ "type": "number", "description": "Page size, 1–250 (default 25)." } - } + }, + "required": [ + "voucherType", + "voucherStatus" + ] }, "endpointMapping": { "method": "GET", diff --git a/packages/backend/src/adapters/intl/api-football.json b/packages/backend/src/adapters/intl/api-football.json index 2fcac4f4..be9124db 100644 --- a/packages/backend/src/adapters/intl/api-football.json +++ b/packages/backend/src/adapters/intl/api-football.json @@ -15,7 +15,7 @@ "name": "API-Football v3", "type": "REST", "baseUrl": "https://v3.football.api-sports.io", - "healthPath": "/status", + "healthcheckPath": "/status", "authType": "API_KEY", "authConfig": { "headerName": "x-apisports-key", diff --git a/packages/backend/src/audit/product-event.service.ts b/packages/backend/src/audit/product-event.service.ts index 0c628b95..b35b34bb 100644 --- a/packages/backend/src/audit/product-event.service.ts +++ b/packages/backend/src/audit/product-event.service.ts @@ -41,6 +41,13 @@ export const ProductEvents = { * brings sign-ups, verified sign-ups and paying customers. */ SIGNUP_ATTRIBUTED: 'signup_attributed', + /** + * A user connected an AI client (Claude, ChatGPT…) through the OAuth flow + * for the first time; metadata.client = the client's name. Server-only. + * Answers: how many sign-ups reach the client, and how many of those then + * add a connector (read against the connectors table). + */ + AI_CLIENT_CONNECTED: 'ai_client_connected', } as const; export type ProductEventName = (typeof ProductEvents)[keyof typeof ProductEvents]; @@ -49,7 +56,10 @@ export type ProductEventName = (typeof ProductEvents)[keyof typeof ProductEvents * Events only the server writes. A signed-in user could otherwise post a * `signup_attributed` of their own and skew the channel report. */ -const SERVER_ONLY = new Set([ProductEvents.SIGNUP_ATTRIBUTED]); +const SERVER_ONLY = new Set([ + ProductEvents.SIGNUP_ATTRIBUTED, + ProductEvents.AI_CLIENT_CONNECTED, +]); const CLIENT_REPORTABLE = new Set( Object.values(ProductEvents).filter((e) => !SERVER_ONLY.has(e)), ); diff --git a/packages/backend/src/common/ssrf.util.spec.ts b/packages/backend/src/common/ssrf.util.spec.ts index 63e7bd1b..de7aef29 100644 --- a/packages/backend/src/common/ssrf.util.spec.ts +++ b/packages/backend/src/common/ssrf.util.spec.ts @@ -14,6 +14,10 @@ describe('extractSsrfBlockedHostname', () => { "SSRF guard: cannot resolve 'other-mcp-server': getaddrinfo ENOTFOUND other-mcp-server", 'other-mcp-server', ], + [ + "Host not found: 'nina.api.proxy.bund.dev' could not be resolved (ENOTFOUND). Check the address in the connector settings.", + 'nina.api.proxy.bund.dev', + ], ])('extracts the host from %s', (message, expected) => { expect(extractSsrfBlockedHostname(message)).toBe(expected); }); @@ -28,3 +32,12 @@ describe('extractSsrfBlockedHostname', () => { expect(extractSsrfBlockedHostname(message)).toBeUndefined(); }); }); + +describe('assertSafeOutboundHost on a name that does not resolve', () => { + it('says the host was not found instead of reporting a policy block', async () => { + const { assertSafeOutboundHost } = await import('./ssrf.util'); + await expect( + assertSafeOutboundHost('no-such-host.invalid', { SSRF_GUARD: 'enabled' } as NodeJS.ProcessEnv), + ).rejects.toThrow(/^Host not found: 'no-such-host\.invalid' could not be resolved \(ENOTFOUND\)/); + }); +}); diff --git a/packages/backend/src/common/ssrf.util.ts b/packages/backend/src/common/ssrf.util.ts index 1e51e289..f8b25c22 100644 --- a/packages/backend/src/common/ssrf.util.ts +++ b/packages/backend/src/common/ssrf.util.ts @@ -243,8 +243,15 @@ async function vetHost( try { resolved = await dns.lookup(hostname, { all: true }); } catch (e: any) { + // Not a policy decision: the name simply has no address (a typo, a + // retired API, a DNS hiccup). Worded as such, because "SSRF guard" made + // users and the model read a security block into a wrong host name. + const temporary = e?.code === 'EAI_AGAIN'; throw new SsrfBlockedError( - `SSRF guard: cannot resolve '${hostname}': ${e?.message || e}`, + `Host not found: '${hostname}' could not be resolved (${e?.code || e?.message || e}). ` + + (temporary + ? 'This is usually a temporary DNS failure; try again.' + : 'Check the address in the connector settings.'), ); } @@ -386,9 +393,9 @@ export function createSsrfGuardedAgents(env: NodeJS.ProcessEnv = process.env): { export function extractSsrfBlockedHostname( message: string, ): string | undefined { - const match = /SSRF guard:\s*(?:address|hostname|cannot resolve)\s*'([^']+)'/.exec( - message || '', - ); + const match = + /SSRF guard:\s*(?:address|hostname|cannot resolve)\s*'([^']+)'/.exec(message || '') ?? + /Host not found:\s*'([^']+)'/.exec(message || ''); return match?.[1]; } diff --git a/packages/backend/src/common/unresolved-placeholders.util.spec.ts b/packages/backend/src/common/unresolved-placeholders.util.spec.ts index 1f6697fc..dc57bc53 100644 --- a/packages/backend/src/common/unresolved-placeholders.util.spec.ts +++ b/packages/backend/src/common/unresolved-placeholders.util.spec.ts @@ -69,6 +69,16 @@ describe('assertNoUnresolvedPlaceholders', () => { ).toThrow(/The "Acme" connector is missing a value for X/); }); + it('links to the connector page when given one', () => { + expect(() => + assertNoUnresolvedPlaceholders( + { authConfig: { token: '{{X}}' } }, + 'the "Acme" connector', + 'https://cloud.example.com/connectors/c1', + ), + ).toThrow(/Open the connector \(https:\/\/cloud\.example\.com\/connectors\/c1\) and set that variable/); + }); + it('falls back to a generic subject', () => { expect(() => assertNoUnresolvedPlaceholders({ authConfig: { token: '{{X}}' } }), diff --git a/packages/backend/src/common/unresolved-placeholders.util.ts b/packages/backend/src/common/unresolved-placeholders.util.ts index 5cd76dfd..2fc6d1a3 100644 --- a/packages/backend/src/common/unresolved-placeholders.util.ts +++ b/packages/backend/src/common/unresolved-placeholders.util.ts @@ -62,6 +62,8 @@ export function assertNoUnresolvedPlaceholders( request: RequestShape, /** How to name the thing in the error, e.g. `the connector behind etsy_get_shop`. */ subject?: string, + /** The connector's page in the dashboard, so the reader can go straight there. */ + fixUrl?: string, ): void { const missing = findUnresolvedPlaceholders({ baseUrl: request.baseUrl, @@ -78,7 +80,7 @@ export function assertNoUnresolvedPlaceholders( `${which} is missing ${missing.length === 1 ? 'a value' : 'values'} for ${names}. ` + 'The request was not sent, because it would have carried the placeholder text ' + 'instead of the credential and the upstream API would have rejected it with a ' + - 'misleading error. Open the connector and set ' + + `misleading error. Open the connector${fixUrl ? ` (${fixUrl})` : ''} and set ` + `${missing.length === 1 ? 'that variable' : 'those variables'}, then try again.`, ); } diff --git a/packages/backend/src/common/url.util.spec.ts b/packages/backend/src/common/url.util.spec.ts index d729e2e6..d1e1a70f 100644 --- a/packages/backend/src/common/url.util.spec.ts +++ b/packages/backend/src/common/url.util.spec.ts @@ -1,5 +1,5 @@ import { BadRequestException } from '@nestjs/common'; -import { normalizeConnectorBaseUrl, resolveMcpEndpointUrl } from './url.util'; +import { connectorPageUrl, normalizeConnectorBaseUrl, resolveMcpEndpointUrl } from './url.util'; describe('normalizeConnectorBaseUrl', () => { it('keeps a well-formed https URL untouched', () => { @@ -155,3 +155,23 @@ describe('resolveMcpEndpointUrl credential handling', () => { ).toBe('http://mcp.example.com:8931/tenant/a'); }); }); + +describe('connectorPageUrl', () => { + const saved = process.env.FRONTEND_URL; + afterEach(() => { + if (saved === undefined) delete process.env.FRONTEND_URL; + else process.env.FRONTEND_URL = saved; + }); + + it('builds the dashboard link from FRONTEND_URL', () => { + process.env.FRONTEND_URL = 'https://cloud.example.com/'; + expect(connectorPageUrl('c1')).toBe('https://cloud.example.com/connectors/c1'); + }); + + it('gives nothing without a usable FRONTEND_URL or id', () => { + delete process.env.FRONTEND_URL; + expect(connectorPageUrl('c1')).toBeUndefined(); + process.env.FRONTEND_URL = 'https://cloud.example.com'; + expect(connectorPageUrl(undefined)).toBeUndefined(); + }); +}); diff --git a/packages/backend/src/common/url.util.ts b/packages/backend/src/common/url.util.ts index 0a6fb14d..2c57673b 100644 --- a/packages/backend/src/common/url.util.ts +++ b/packages/backend/src/common/url.util.ts @@ -144,3 +144,15 @@ function withPath(base: URL, pathname: string, search: string): URL { url.hash = ''; return url; } + +/** + * Public dashboard address of a connector's page, for messages a person reads + * in a chat client ("open the connector and set X"). FRONTEND_URL is where the + * dashboard is served; undefined when it is not configured, so callers can + * fall back to wording without a link. + */ +export function connectorPageUrl(connectorId: string | undefined | null): string | undefined { + const base = (process.env.FRONTEND_URL || '').trim().replace(/\/+$/, ''); + if (!connectorId || !/^https?:\/\//i.test(base)) return undefined; + return `${base}/connectors/${encodeURIComponent(connectorId)}`; +} diff --git a/packages/backend/src/connectors/connectors.service.ts b/packages/backend/src/connectors/connectors.service.ts index e63f3f1c..6c96fed6 100644 --- a/packages/backend/src/connectors/connectors.service.ts +++ b/packages/backend/src/connectors/connectors.service.ts @@ -18,7 +18,7 @@ import { CALLER_CONTEXT_PREFIX } from '../common/caller-context.util'; import { assertNoUnresolvedPlaceholders } from '../common/unresolved-placeholders.util'; import { assertAbsoluteBaseUrl } from '../common/base-url-variable.util'; import { extractSsrfBlockedHostname } from '../common/ssrf.util'; -import { normalizeConnectorBaseUrl } from '../common/url.util'; +import { connectorPageUrl, normalizeConnectorBaseUrl } from '../common/url.util'; import { resolveAdapterIcon } from './connector-icon.util'; import { applySchemaDefaults } from '../common/schema-defaults.util'; import { renderStaticResponse } from './static-response.util'; @@ -251,6 +251,7 @@ export class ConnectorsService { assertNoUnresolvedPlaceholders( { baseUrl, headers, authConfig }, `the "${connector.name}" connector`, + connectorPageUrl(connector.id), ); assertAbsoluteBaseUrl( { @@ -547,6 +548,7 @@ export class ConnectorsService { authConfig, }, toolName ? `the connector behind ${toolName}` : `the "${connector.name}" connector`, + connectorPageUrl(connector.id), ); assertAbsoluteBaseUrl( { diff --git a/packages/backend/src/connectors/engines/oauth2-token.service.ts b/packages/backend/src/connectors/engines/oauth2-token.service.ts index ed169524..2a0129f2 100644 --- a/packages/backend/src/connectors/engines/oauth2-token.service.ts +++ b/packages/backend/src/connectors/engines/oauth2-token.service.ts @@ -12,6 +12,7 @@ import { isPrivateKeyJwt, } from './client-assertion.util'; import { ssrfGuardedAxiosOptions } from '../../common/guarded-http.util'; +import { connectorPageUrl } from '../../common/url.util'; /** Refresh tokens that expire within this window (5 minutes). */ const PROACTIVE_REFRESH_BUFFER_MS = 5 * 60 * 1000; @@ -147,9 +148,10 @@ export class OAuth2TokenService { ); } if (!authConfig.refreshToken && authConfig.authorizationUrl) { + const page = connectorPageUrl(connectorId); throw unauthorized( 'OAuth2: this connector has not been authorized yet. No request was sent to the API. ' + - 'Open the connector in AnythingMCP and click Authorize with Provider.', + `Open the connector in AnythingMCP${page ? ` (${page})` : ''} and click Authorize with Provider.`, ); } } diff --git a/packages/backend/src/connectors/engines/rest.engine.spec.ts b/packages/backend/src/connectors/engines/rest.engine.spec.ts index d03a1f28..6647e753 100644 --- a/packages/backend/src/connectors/engines/rest.engine.spec.ts +++ b/packages/backend/src/connectors/engines/rest.engine.spec.ts @@ -1,4 +1,4 @@ -import { RestEngine, serializeRepeatedParams } from './rest.engine'; +import { RestEngine, parseRetryAfterMs, serializeRepeatedParams } from './rest.engine'; import { OAuth2TokenService } from './oauth2-token.service'; import { LoginTokenService } from './login-token.service'; import axios, { AxiosError } from 'axios'; @@ -743,6 +743,49 @@ describe('RestEngine', () => { expect(mockedAxios).toHaveBeenCalledTimes(1); }); + // A rate limit is not an outage: at most one more attempt, and only when + // the API's Retry-After fits inside a tool call. + describe('rate limits', () => { + const limited = (status: number, retryAfter?: string) => + new AxiosError('limited', undefined, undefined, {}, { + status, + data: {}, + headers: retryAfter === undefined ? {} : { 'retry-after': retryAfter }, + } as any); + const call = () => + engine.execute( + { baseUrl: 'https://api.example.com', authType: 'NONE' }, + { method: 'GET', path: '/' }, + {}, + ); + + it('tries a 429 without Retry-After only once more', async () => { + mockedAxios.mockRejectedValue(limited(429)); + await expect(call()).rejects.toBeInstanceOf(AxiosError); + expect(mockedAxios).toHaveBeenCalledTimes(2); + }); + + it('honours a short Retry-After and returns the success', async () => { + mockedAxios + .mockRejectedValueOnce(limited(429, '0')) + .mockResolvedValueOnce({ data: { ok: true } }); + await expect(call()).resolves.toEqual({ ok: true }); + expect(mockedAxios).toHaveBeenCalledTimes(2); + }); + + it('does not retry when Retry-After is longer than a tool call can wait', async () => { + mockedAxios.mockRejectedValue(limited(429, '60')); + await expect(call()).rejects.toBeInstanceOf(AxiosError); + expect(mockedAxios).toHaveBeenCalledTimes(1); + }); + + it('treats a 503 with a long Retry-After the same way', async () => { + mockedAxios.mockRejectedValue(limited(503, '120')); + await expect(call()).rejects.toBeInstanceOf(AxiosError); + expect(mockedAxios).toHaveBeenCalledTimes(1); + }); + }); + it('gives up after exhausting retries on persistent 503', async () => { mockedAxios.mockRejectedValue(err(503)); @@ -1460,3 +1503,20 @@ describe('RestEngine — bodyTemplate that will not parse', () => { expect(err.message).not.toMatch(/super-secret/); }); }); + +describe('parseRetryAfterMs', () => { + it('reads delay-seconds', () => { + expect(parseRetryAfterMs('2')).toBe(2000); + expect(parseRetryAfterMs(['5'])).toBe(5000); + }); + it('reads an HTTP date relative to now', () => { + const now = Date.parse('Sat, 03 Oct 2026 10:00:00 GMT'); + expect(parseRetryAfterMs('Sat, 03 Oct 2026 10:00:02 GMT', now)).toBe(2000); + expect(parseRetryAfterMs('Sat, 03 Oct 2026 09:00:00 GMT', now)).toBe(0); + }); + it('returns null for missing or unreadable values', () => { + expect(parseRetryAfterMs(undefined)).toBeNull(); + expect(parseRetryAfterMs('')).toBeNull(); + expect(parseRetryAfterMs('soon')).toBeNull(); + }); +}); diff --git a/packages/backend/src/connectors/engines/rest.engine.ts b/packages/backend/src/connectors/engines/rest.engine.ts index 57412fea..486bf05c 100644 --- a/packages/backend/src/connectors/engines/rest.engine.ts +++ b/packages/backend/src/connectors/engines/rest.engine.ts @@ -26,6 +26,23 @@ import { ssrfGuardedAxiosOptions } from '../../common/guarded-http.util'; * Supports OAuth2 token refresh: if a request returns 401 and a refreshToken + tokenUrl * are available, it will attempt to refresh the access token and retry the request once. */ +/** Longest Retry-After we honour inside a tool call; beyond it we give up. */ +const MAX_RETRY_AFTER_MS = 3000; + +/** + * Retry-After as milliseconds: delay-seconds or an HTTP date (RFC 9110 + * 10.2.3). Null when absent or unreadable. + */ +export function parseRetryAfterMs(value: unknown, now = Date.now()): number | null { + if (value === undefined || value === null) return null; + const text = String(Array.isArray(value) ? value[0] : value).trim(); + if (!text) return null; + if (/^\d+$/.test(text)) return Number(text) * 1000; + const at = Date.parse(text); + if (Number.isNaN(at)) return null; + return Math.max(0, at - now); +} + @Injectable() export class RestEngine { private readonly logger = new Logger(RestEngine.name); @@ -388,6 +405,34 @@ export class RestEngine { ].includes(error.code ?? ''); } + /** + * How long to wait before the next attempt, or null to give up. + * + * A rate limit is not an outage: retrying a 429 three times in 3.7 s sent + * four requests to an API that had just asked for fewer, and a customer + * whose backend hit its own Supabase limit saw ~54k rejections become ~216k + * requests in two days. So a 429, and a 503 that says when to come back, + * get at most one more attempt: after the API's Retry-After when it is + * short, after about a second when it gives none, and none at all when it + * asks for longer than we can wait inside a tool call. The model still sees + * the status and the Retry-After header in the error. + */ + private retryDelayMs( + error: unknown, + attempt: number, + delaysMs: number[], + ): number | null { + const response = (error as AxiosError).response; + const status = response?.status; + const retryAfter = parseRetryAfterMs(response?.headers?.['retry-after']); + if (status === 429 || (status === 503 && retryAfter !== null)) { + if (attempt > 0) return null; + if (retryAfter === null) return 1000 + Math.floor(Math.random() * 250); + return retryAfter <= MAX_RETRY_AFTER_MS ? retryAfter : null; + } + return attempt < delaysMs.length ? delaysMs[attempt] : null; + } + /** * Human-readable replacements for connection-level failures. Without these * the caller — and the model reading the tool result — gets the raw OpenSSL @@ -441,7 +486,8 @@ export class RestEngine { return await axios(axiosConfig); } catch (error) { const transient = this.isTransientError(error); - if (attempt >= delaysMs.length || !transient) { + const delay = transient ? this.retryDelayMs(error, attempt, delaysMs) : null; + if (delay === null) { // Warn, not debug: production runs at info, so a call that burned // every retry used to look identical to one that failed outright — // there was no way to tell a flaky upstream from a broken connector. @@ -456,9 +502,9 @@ export class RestEngine { throw this.describeConnectionError(error, attempt + 1); } this.logger.debug( - `Transient error (attempt ${attempt + 1}), retrying in ${delaysMs[attempt]}ms`, + `Transient error (attempt ${attempt + 1}), retrying in ${delay}ms`, ); - await new Promise((resolve) => setTimeout(resolve, delaysMs[attempt])); + await new Promise((resolve) => setTimeout(resolve, delay)); } } } diff --git a/packages/backend/src/ee/cloud/onboarding-cron.service.spec.ts b/packages/backend/src/ee/cloud/onboarding-cron.service.spec.ts index c3d27b18..b7771753 100644 --- a/packages/backend/src/ee/cloud/onboarding-cron.service.spec.ts +++ b/packages/backend/src/ee/cloud/onboarding-cron.service.spec.ts @@ -202,3 +202,98 @@ describe('OnboardingCronService — trial repair', () => { expect(out.licensesDeactivated).toBe(1); }); }); + +describe('OnboardingCronService — onboarding pass', () => { + const HOUR = 60 * 60 * 1000; + function makeService(opts: { + candidates: any[]; + orgConnectors?: Record; + grants?: { userId: string; clientId: string }[]; + }) { + const update = jest.fn().mockResolvedValue({}); + const prisma = { + user: { + findMany: jest + .fn() + .mockResolvedValueOnce(opts.candidates) + .mockResolvedValueOnce([]), + update, + }, + connector: { + groupBy: jest.fn().mockResolvedValue( + Object.entries(opts.orgConnectors ?? {}).map(([organizationId, n]) => ({ + organizationId, + _count: { _all: n }, + })), + ), + }, + mcpConnectionGrant: { findMany: jest.fn().mockResolvedValue(opts.grants ?? []) }, + oAuthClient: { + findMany: jest.fn().mockResolvedValue([{ clientId: 'claude-client', clientName: 'Claude' }]), + }, + license: { + findMany: jest.fn().mockResolvedValue([]), + updateMany: jest.fn().mockResolvedValue({ count: 0 }), + }, + } as any; + const email = { + sendOnboardingReminderEmail: jest.fn().mockResolvedValue(true), + sendActivationReminderEmail: jest.fn().mockResolvedValue(true), + } as any; + return { service: new OnboardingCronService(prisma, email, makeLicense()), email, update, prisma }; + } + const user = (over: Partial) => ({ + id: 'u1', + email: 'u1@example.com', + name: 'Ada', + organizationId: 'org-1', + onboardingCompletedAt: null, + onboardingReminderCount: 0, + onboardingLastReminderAt: null, + _count: { connectors: 0 }, + createdAt: new Date(Date.now() - 3 * HOUR), + ...over, + }); + + it('nudges a user who connected Claude to an empty workspace, naming the client', async () => { + const { service, email, update } = makeService({ + candidates: [user({})], + grants: [{ userId: 'u1', clientId: 'claude-client' }], + }); + await service.run(); + expect(email.sendOnboardingReminderEmail).toHaveBeenCalledWith('u1@example.com', 'Ada', 1, { + aiClient: 'Claude', + }); + expect(update).toHaveBeenCalledWith( + expect.objectContaining({ data: expect.objectContaining({ onboardingReminderCount: 1 }) }), + ); + }); + + it('waits the usual 24h for a user with no AI client yet', async () => { + const { service, email } = makeService({ candidates: [user({})] }); + await service.run(); + expect(email.sendOnboardingReminderEmail).not.toHaveBeenCalled(); + }); + + it('still reminds a user who pressed Skip on /welcome with an empty workspace', async () => { + const { service, email } = makeService({ + candidates: [ + user({ onboardingCompletedAt: new Date(), createdAt: new Date(Date.now() - 30 * HOUR) }), + ], + }); + await service.run(); + expect(email.sendOnboardingReminderEmail).toHaveBeenCalledWith('u1@example.com', 'Ada', 1, undefined); + }); + + it('counts a teammate\'s connector: no reminder, completion stamped', async () => { + const { service, email, update } = makeService({ + candidates: [user({ createdAt: new Date(Date.now() - 30 * HOUR) })], + orgConnectors: { 'org-1': 1 }, + }); + await service.run(); + expect(email.sendOnboardingReminderEmail).not.toHaveBeenCalled(); + expect(update).toHaveBeenCalledWith( + expect.objectContaining({ data: { onboardingCompletedAt: expect.any(Date) } }), + ); + }); +}); diff --git a/packages/backend/src/ee/cloud/onboarding-cron.service.ts b/packages/backend/src/ee/cloud/onboarding-cron.service.ts index e8d5390c..59c68b37 100644 --- a/packages/backend/src/ee/cloud/onboarding-cron.service.ts +++ b/packages/backend/src/ee/cloud/onboarding-cron.service.ts @@ -4,6 +4,9 @@ import { EmailService } from '../../settings/email.service'; import { LicenseService } from '../../license/license.service'; const HOURS = (n: number) => n * 60 * 60 * 1000; + +/** How long after connecting an AI client to an empty workspace we nudge. */ +const AI_CLIENT_NUDGE_AFTER = HOURS(2); const DAYS = (n: number) => n * 24 * 60 * 60 * 1000; /** @@ -67,18 +70,23 @@ export class OnboardingCronService { skipped: 0, }; - // Candidate set: verified, no completion, ≤2 reminders, not opted out, - // and registered between 24h and 14d ago. We bound at 14d so a user + // Candidate set: verified, ≤2 reminders, not opted out, registered + // between AI_CLIENT_NUDGE_AFTER and 14d ago. We bound at 14d so a user // who signed up months ago doesn't suddenly get woken up if we ever // backfill columns. + // + // onboardingCompletedAt is NOT a filter any more: the Skip button on + // /welcome sets it, and people who skipped the page with an empty + // workspace stopped getting the reminders that were meant for exactly + // them (13 of the 428 sign-ups of 1-3 Oct 2026). What counts is whether + // the workspace has a connector, checked below. const candidates = await this.prisma.user.findMany({ where: { emailVerified: true, emailMarketingOptOut: false, - onboardingCompletedAt: null, onboardingReminderCount: { lt: 2 }, createdAt: { - lte: new Date(now - HOURS(24)), + lte: new Date(now - AI_CLIENT_NUDGE_AFTER), gte: new Date(now - HOURS(24 * 14)), }, }, @@ -87,25 +95,56 @@ export class OnboardingCronService { email: true, name: true, createdAt: true, + organizationId: true, + onboardingCompletedAt: true, onboardingReminderCount: true, onboardingLastReminderAt: true, _count: { select: { connectors: true } }, }, }); + // Connectors per workspace: a teammate's connector serves this user too. + const orgIds = [ + ...new Set(candidates.map((u) => u.organizationId).filter((id): id is string => !!id)), + ]; + const orgConnectors = new Map(); + if (orgIds.length > 0) { + const rows = await this.prisma.connector.groupBy({ + by: ['organizationId'], + where: { organizationId: { in: orgIds } }, + _count: { _all: true }, + }); + for (const r of rows) { + if (r.organizationId) orgConnectors.set(r.organizationId, r._count._all); + } + } + + // Who already connected an AI client (Claude, ChatGPT…) and how long ago. + // These are the warmest leads of all: the client is waiting on a + // workspace with nothing in it. + const aiClients = await this.connectedAiClients( + candidates.map((u) => u.id), + now, + ); + for (const u of candidates) { out.examined++; - // Race-safe: a user that created a connector between candidate - // pull and now should never receive a nudge. - if (u._count.connectors > 0) { + // Race-safe: a workspace that got a connector between candidate pull + // and now should never receive a nudge. + const connectors = + (u.organizationId ? orgConnectors.get(u.organizationId) : undefined) ?? + u._count.connectors; + if (connectors > 0) { // Auto-stamp completion so we never see them again. - await this.prisma.user - .update({ - where: { id: u.id }, - data: { onboardingCompletedAt: new Date() }, - }) - .catch(() => {}); + if (!u.onboardingCompletedAt) { + await this.prisma.user + .update({ + where: { id: u.id }, + data: { onboardingCompletedAt: new Date() }, + }) + .catch(() => {}); + } out.skipped++; continue; } @@ -115,12 +154,16 @@ export class OnboardingCronService { ? now - u.onboardingLastReminderAt.getTime() : Infinity; - // First nudge: 24-72h after signup, count == 0. - if (u.onboardingReminderCount === 0 && age >= HOURS(24)) { + // First nudge: 24h after signup, or as soon as an AI client has been + // connected for AI_CLIENT_NUDGE_AFTER, whichever comes first. The + // second one names the client, because that is what the user just did. + const aiClient = aiClients.get(u.id); + if (u.onboardingReminderCount === 0 && (age >= HOURS(24) || aiClient)) { const ok = await this.email.sendOnboardingReminderEmail( u.email, u.name || 'there', 1, + aiClient ? { aiClient } : undefined, ); if (ok) { await this.prisma.user.update({ @@ -310,6 +353,33 @@ export class OnboardingCronService { * `firstSuccessfulInvocationAt: null` = never activated; `activationReminderAt` * caps it at one send. */ + /** + * Users among `userIds` with a live AI-client connection made at least + * AI_CLIENT_NUDGE_AFTER ago, mapped to the client's display name. + */ + private async connectedAiClients(userIds: string[], now: number): Promise> { + const out = new Map(); + if (userIds.length === 0) return out; + const grants = await this.prisma.mcpConnectionGrant.findMany({ + where: { + userId: { in: userIds }, + revokedAt: null, + createdAt: { lte: new Date(now - AI_CLIENT_NUDGE_AFTER) }, + }, + select: { userId: true, clientId: true }, + }); + if (grants.length === 0) return out; + const clients = await this.prisma.oAuthClient.findMany({ + where: { clientId: { in: [...new Set(grants.map((g) => g.clientId))] } }, + select: { clientId: true, clientName: true }, + }); + const names = new Map(clients.map((c) => [c.clientId, c.clientName])); + for (const g of grants) { + if (!out.has(g.userId)) out.set(g.userId, names.get(g.clientId) || 'your AI client'); + } + return out; + } + private async runActivationPass( now: number, out: { diff --git a/packages/backend/src/mcp-server/dynamic-mcp-tools.ts b/packages/backend/src/mcp-server/dynamic-mcp-tools.ts index 611f1ca4..6cffc921 100644 --- a/packages/backend/src/mcp-server/dynamic-mcp-tools.ts +++ b/packages/backend/src/mcp-server/dynamic-mcp-tools.ts @@ -35,6 +35,7 @@ import { processGauges } from '../common/process-vitals'; import { applySchemaDefaults } from '../common/schema-defaults.util'; import { renderStaticResponse } from '../connectors/static-response.util'; import { ODataEngine, isODataBuiltinMethod } from '../connectors/engines/odata.engine'; +import { connectorPageUrl } from '../common/url.util'; /** * ToolExecutor — executes dynamically registered MCP tools. @@ -351,6 +352,7 @@ export class DynamicMcpTools { authConfig, }, `the connector behind ${tool.name}`, + connectorPageUrl(tool.connectorId), ); // A base URL without https:// (a variable typed as `shop.example.com` // on an older install) would otherwise fail in the SSRF guard as diff --git a/packages/backend/src/mcp-server/error-hints.spec.ts b/packages/backend/src/mcp-server/error-hints.spec.ts index 0127a720..580b44ba 100644 --- a/packages/backend/src/mcp-server/error-hints.spec.ts +++ b/packages/backend/src/mcp-server/error-hints.spec.ts @@ -127,4 +127,34 @@ describe('deriveErrorHint — SQL-backed customer APIs', () => { it('still says nothing about an error it does not recognise', () => { expect(deriveErrorHint({ status: 500, body: { message: 'boom' } })).toBeUndefined(); }); + + describe('Telegram', () => { + const host = 'api.telegram.org'; + it('points at the bot token on a 404 Not Found', () => { + expect( + deriveErrorHint({ host, status: 404, body: { ok: false, error_code: 404, description: 'Not Found' } }), + ).toMatch(/TELEGRAM_BOT_TOKEN/); + }); + it('explains how a bot reaches a chat', () => { + expect( + deriveErrorHint({ host, status: 400, body: { description: 'Bad Request: chat not found' } }), + ).toMatch(/pressed Start/); + expect( + deriveErrorHint({ host, status: 403, body: { description: "Forbidden: bot can't initiate conversation with a user" } }), + ).toMatch(/pressed Start/); + }); + it('stays quiet for another host answering Not Found', () => { + expect(deriveErrorHint({ host: 'api.example.com', status: 404, body: '"Not Found"' })).toBeUndefined(); + }); + }); + + it('tells the model to split Lexware overdue from other statuses', () => { + expect( + deriveErrorHint({ + host: 'api.lexware.io', + status: 400, + body: { message: "voucherStatus filter 'overdue' cannot be used in combination with other states" }, + }), + ).toMatch(/one call with voucherStatus=overdue/); + }); }); diff --git a/packages/backend/src/mcp-server/error-hints.ts b/packages/backend/src/mcp-server/error-hints.ts index 154fdb1b..da61a8b1 100644 --- a/packages/backend/src/mcp-server/error-hints.ts +++ b/packages/backend/src/mcp-server/error-hints.ts @@ -87,6 +87,24 @@ const TYPESAFE_INVALID_REQUEST_HINT = 'state, model and questions go in the request: there is no field for a list ' + 'of records, send one call per record instead.'; +const TELEGRAM_BAD_TOKEN_HINT = + 'Telegram answers 404 "Not Found" (or 401) for every method when the bot token ' + + 'in the URL is wrong, so this is the TELEGRAM_BOT_TOKEN, not the method. Tell ' + + 'the user to copy the token again from @BotFather (format 123456789:AA...) into ' + + 'the connector settings. Retrying or calling another method will fail the same way.'; + +const TELEGRAM_CHAT_HINT = + 'The bot cannot reach that chat. A bot may only write to a user who has opened ' + + 'it and pressed Start, or to a group/channel it was added to (channels: as an ' + + 'administrator). Use the numeric chat id from telegram_bot_get_updates after ' + + 'the user has sent the bot a message; @usernames only work for public channels. ' + + 'Ask the user to do that instead of trying other ids.'; + +const LEXWARE_OVERDUE_HINT = + 'Lexware does not accept `overdue` together with other statuses in ' + + '`voucherStatus`. Make one call with voucherStatus=overdue and a separate one ' + + 'for the other statuses.'; + function bodyText(body: unknown): string { if (body === undefined || body === null) return ''; if (typeof body === 'string') return body; @@ -142,6 +160,19 @@ export function deriveErrorHint(input: ErrorHintInput): string | undefined { } } + if (hostMatches(input.host, 'api.telegram.org')) { + if (input.status === 401 || (input.status === 404 && /"Not Found"/.test(text))) { + return TELEGRAM_BAD_TOKEN_HINT; + } + if (/chat not found|bot is not a member|can't initiate conversation|need administrator rights|bot was blocked by the user/i.test(text)) { + return TELEGRAM_CHAT_HINT; + } + } + + if (/voucherStatus filter 'overdue' cannot be used in combination/i.test(text)) { + return LEXWARE_OVERDUE_HINT; + } + // TypeSafe answers a question with an unknown type (and any other field it // does not expect) with a bare "Invalid request.", which names nothing the // model could fix. diff --git a/packages/backend/src/mcp-servers/mcp-connection-grant.service.spec.ts b/packages/backend/src/mcp-servers/mcp-connection-grant.service.spec.ts index 71996b13..7572f611 100644 --- a/packages/backend/src/mcp-servers/mcp-connection-grant.service.spec.ts +++ b/packages/backend/src/mcp-servers/mcp-connection-grant.service.spec.ts @@ -346,3 +346,36 @@ describe('McpConnectionGrantService writing a grant', () => { expect(calls.upsert[0].update.revokedAt).toBeNull(); }); }); + +describe('McpConnectionGrantService reports the first connection of a client', () => { + function withEvents(existing: object | null) { + const { prisma } = build({ + grant: existing as any, + servers: [{ id: 'srv-1', organizationId: 'org-1' }], + memberCount: 1, + }); + prisma.oAuthClient = { + findUnique: jest.fn().mockResolvedValue({ clientName: 'Claude' }), + }; + const events = { log: jest.fn().mockResolvedValue(undefined) }; + const svc = new McpConnectionGrantService(prisma, events as any); + return { svc, events }; + } + + it('logs ai_client_connected with the client name on a new grant', async () => { + const { svc, events } = withEvents(null); + await svc.grantServers('client-1', 'user-1', ['srv-1']); + expect(events.log).toHaveBeenCalledWith({ + event: 'ai_client_connected', + userId: 'user-1', + organizationId: 'org-1', + metadata: { client: 'Claude' }, + }); + }); + + it('does not log again when the user changes an existing grant', async () => { + const { svc, events } = withEvents({ id: 'g1', organizationId: null, serverIds: ['srv-1'] }); + await svc.grantWholeOrganization('client-1', 'user-1', 'org-1'); + expect(events.log).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/backend/src/mcp-servers/mcp-connection-grant.service.ts b/packages/backend/src/mcp-servers/mcp-connection-grant.service.ts index fdd1c747..4bf2347e 100644 --- a/packages/backend/src/mcp-servers/mcp-connection-grant.service.ts +++ b/packages/backend/src/mcp-servers/mcp-connection-grant.service.ts @@ -1,5 +1,6 @@ -import { Injectable, Logger } from '@nestjs/common'; +import { Injectable, Logger, Optional } from '@nestjs/common'; import { PrismaService } from '../common/prisma.service'; +import { ProductEvents, ProductEventService } from '../audit/product-event.service'; /** * What a grant resolved to, for the caller that is about to build a tool list. @@ -42,7 +43,10 @@ export type ResolvedGrant = export class McpConnectionGrantService { private readonly logger = new Logger(McpConnectionGrantService.name); - constructor(private readonly prisma: PrismaService) {} + constructor( + private readonly prisma: PrismaService, + @Optional() private readonly events?: ProductEventService, + ) {} /** * Resolve what `clientId` may see on behalf of `userId`, re-validating every @@ -93,7 +97,7 @@ export class McpConnectionGrantService { ); return false; } - await this.upsert(clientId, userId, { organizationId, serverIds: [] }); + await this.upsert(clientId, userId, { organizationId, serverIds: [] }, organizationId); return true; } @@ -118,10 +122,12 @@ export class McpConnectionGrantService { } if (valid.length === 0) return []; - await this.upsert(clientId, userId, { - organizationId: null, - serverIds: valid.map((s) => s.id), - }); + await this.upsert( + clientId, + userId, + { organizationId: null, serverIds: valid.map((s) => s.id) }, + valid[0].organizationId, + ); return valid.map((s) => s.id); } @@ -251,7 +257,13 @@ export class McpConnectionGrantService { clientId: string, userId: string, data: { organizationId: string | null; serverIds: string[] }, + /** Workspace the grant points into, for the product event. */ + eventOrganizationId?: string, ): Promise { + const existing = await this.prisma.mcpConnectionGrant.findUnique({ + where: { clientId_userId: { clientId, userId } }, + select: { id: true }, + }); await this.prisma.mcpConnectionGrant.upsert({ where: { clientId_userId: { clientId, userId } }, create: { clientId, userId, ...data }, @@ -259,5 +271,24 @@ export class McpConnectionGrantService { // client may reach, which is the opposite of having revoked it. update: { ...data, revokedAt: null }, }); + if (!existing) await this.reportConnected(clientId, userId, eventOrganizationId); + } + + /** First connection of this client for this user: the funnel step between sign-up and first call. */ + private async reportConnected( + clientId: string, + userId: string, + organizationId?: string, + ): Promise { + if (!this.events) return; + const client = await this.prisma.oAuthClient + .findUnique({ where: { clientId }, select: { clientName: true } }) + .catch(() => null); + await this.events.log({ + event: ProductEvents.AI_CLIENT_CONNECTED, + userId, + organizationId: organizationId ?? null, + metadata: { client: client?.clientName ?? 'unknown' }, + }); } } diff --git a/packages/backend/src/settings/email.service.ts b/packages/backend/src/settings/email.service.ts index c1f113bd..e9bd1be9 100644 --- a/packages/backend/src/settings/email.service.ts +++ b/packages/backend/src/settings/email.service.ts @@ -6,6 +6,16 @@ import { OrgSettingsService } from './org-settings.service'; import { PrismaService } from '../common/prisma.service'; import { DeploymentService } from '../common/deployment.service'; +/** Escape text placed into an email's HTML (names and client names are user-chosen). */ +function escapeHtml(value: string): string { + return value + .replace(/&/g, '&') + .replace(//g, '>') + .replace(/"/g, '"') + .replace(/'/g, '''); +} + // Production always talks to anythingmcp.com. The licence site decides // which plan an installation runs; a URL taken from the environment would let // any self-hosted operator point verification at a server of their own and @@ -471,6 +481,8 @@ export class EmailService { to: string, name: string, dayNumber: 1 | 2, + /** Set when the user already connected an AI client to the empty workspace. */ + opts?: { aiClient?: string }, ): Promise { const transport = await this.createTransporter(); if (!transport) { @@ -483,20 +495,28 @@ export class EmailService { const cloudUrl = process.env.CLOUD_PUBLIC_URL || 'https://cloud.anythingmcp.com'; const welcomeUrl = `${cloudUrl}/welcome`; + const storeUrl = `${cloudUrl}/connectors/store`; const unsubUrl = `${cloudUrl}/settings/profile`; + const client = opts?.aiClient ? escapeHtml(opts.aiClient) : undefined; + const safeName = escapeHtml(name); - const subject = - dayNumber === 1 + const subject = client + ? `${opts!.aiClient} is connected. Now give it something to work with` + : dayNumber === 1 ? 'Connect your first tool in 60 seconds — AnythingMCP' : 'Still here? Pick a tool to try — AnythingMCP'; - const body = - dayNumber === 1 - ? `

Hi ${name},

-

You signed up for AnythingMCP yesterday but haven't connected anything yet. The fastest path to your first AI superpower is picking a ready-made connector from the marketplace — Sendcloud, Stripe, GitHub, Slack, Help Scout… 180+ are pre-wired.

+ const body = client + ? `

Hi ${safeName},

+

You connected ${client} to AnythingMCP, but your workspace has no connectors yet, so ${client} has nothing to reach.

+

Add the app you want it to work with: Etsy, Odoo, weclapp, Lexware, Telegram, Shopify and 260 more are ready to install. As soon as one is in, ask ${client} about it in the same chat.

+

Add your first connector →

` + : dayNumber === 1 + ? `

Hi ${safeName},

+

You signed up for AnythingMCP yesterday but haven't connected anything yet. The fastest path to your first AI superpower is picking a ready-made connector from the marketplace: Etsy, Odoo, weclapp, Lexware, Sendcloud, GitHub… 265 are pre-wired.

Open the welcome wizard →

Should take about a minute.

` - : `

Hi ${name},

+ : `

Hi ${safeName},

Just checking in — your AnythingMCP account is still waiting for its first connector. If anything got in your way, hit reply and tell us what; we read every reply.

Pick a connector →

`; @@ -516,13 +536,17 @@ export class EmailService { `, text: `Hi ${name},\n\n${ - dayNumber === 1 - ? "You signed up for AnythingMCP yesterday but haven't connected anything yet." - : 'Your AnythingMCP account is still waiting for its first connector.' - }\n\nOpen the wizard: ${welcomeUrl}\n\nUnsubscribe: ${unsubUrl}`, + opts?.aiClient + ? `You connected ${opts.aiClient} to AnythingMCP, but your workspace has no connectors yet.\n\nAdd your first connector: ${storeUrl}` + : `${ + dayNumber === 1 + ? "You signed up for AnythingMCP yesterday but haven't connected anything yet." + : 'Your AnythingMCP account is still waiting for its first connector.' + }\n\nOpen the wizard: ${welcomeUrl}` + }\n\nUnsubscribe: ${unsubUrl}`, }); this.logger.log( - `Onboarding-reminder email (day ${dayNumber}) sent to ${to}`, + `Onboarding-reminder email (day ${dayNumber}${opts?.aiClient ? ', AI client connected' : ''}) sent to ${to}`, ); return true; } catch (err) {