Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
9f81f6f
fix(serializer): skip tool selection for trigger-mode blocks
waleedlatif1 Oct 9, 2026
5d792d0
improvement(executor): memoize block output schemas per serialized block
waleedlatif1 Oct 9, 2026
cf9f356
fix(webhooks): stop warning on accountless webhook credentials
waleedlatif1 Oct 9, 2026
b1257f4
improvement(webhooks): log Slack reactions.get missing_scope once per…
waleedlatif1 Oct 9, 2026
0e55374
improvement(workflows): repair drifted subBlock types quietly
waleedlatif1 Oct 9, 2026
9862356
improvement(catalog): skip custom block input derivation for tool reads
waleedlatif1 Oct 9, 2026
382dbc1
improvement(webhooks): debug-log missing idempotency ids for provider…
waleedlatif1 Oct 9, 2026
c349281
fix(flint): declare generate_pages items as an array param
waleedlatif1 Oct 9, 2026
2f2bfe8
improvement(executor): skip JSON parsing for json inputs holding plai…
waleedlatif1 Oct 9, 2026
5137ea5
improvement(mothership): classify catalog routes and dedupe missing-p…
waleedlatif1 Oct 9, 2026
11c2c36
fix(tokenization): resolve tiktoken encodings by model family
waleedlatif1 Oct 9, 2026
f6e2c8b
improvement(proxy): log blocked scanner requests at debug
waleedlatif1 Oct 9, 2026
7a40330
improvement(telemetry): downgrade best-effort collector failures to warn
waleedlatif1 Oct 9, 2026
17b4146
fix(admin): log mothership admin proxy failures and map upstream 5xx …
waleedlatif1 Oct 9, 2026
d1731a4
improvement(logs): move per-call success chatter to debug
waleedlatif1 Oct 9, 2026
5d32640
chore(audits): shrink explicit-any baseline after tokenization cleanup
waleedlatif1 Oct 10, 2026
285b846
improvement(mothership): keep tool call routing decision at info
waleedlatif1 Oct 10, 2026
78aec62
fix(admin): keep upstream message and empty 2xx bodies in the mothers…
waleedlatif1 Oct 10, 2026
683a7aa
improvement(tokenization): memoize model to encoding name resolution
waleedlatif1 Oct 10, 2026
965ddea
improvement(mothership): resolve catalog route patterns through the r…
waleedlatif1 Oct 10, 2026
b5e0932
fix(webhooks): warn on missing idempotency ids unless the provider de…
waleedlatif1 Oct 10, 2026
8fc6dc5
refactor(logs): simplify log-noise fixes after review
waleedlatif1 Oct 10, 2026
52adaea
chore(logs): tighten comments and rename trigger-mode serializer test
waleedlatif1 Oct 10, 2026
d08916b
test(mothership): guard sandbox catalog route classification against …
waleedlatif1 Oct 10, 2026
e4659a4
fix(webhooks): keep calling reactions.get after a missing-scope report
waleedlatif1 Oct 10, 2026
39adc46
chore(logs): tighten log-noise changes after a line-by-line audit
waleedlatif1 Oct 10, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/docs/content/docs/integrations/flint.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ Start a background Flint agent task that generates up to 10 pages from a templat
| `apiKey` | string | Yes | Flint API key \(found in Flint team settings, starts with ak_\) |
| `siteId` | string | Yes | ID of the Flint site the agent should modify |
| `templatePageSlug` | string | Yes | Slug of the existing template page to generate from \(e.g., /case-studies/template\) |
| `items` | json | Yes | JSON array of 1-10 pages to generate. Each item requires targetPageSlug \(slug for the new page\) and context \(content details the agent should use\). |
| `items` | array | Yes | JSON array of 1-10 pages to generate. Each item requires targetPageSlug \(slug for the new page\) and context \(content details the agent should use\). |
| `callbackUrl` | string | No | HTTPS webhook URL that Flint will POST to when the task completes or fails |
| `publish` | boolean | No | Whether to automatically publish the generated pages when the task completes |

Expand Down
214 changes: 95 additions & 119 deletions apps/sim/app/api/admin/mothership/route.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
import { db } from '@sim/db'
import { settings, user } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { getErrorMessage } from '@sim/utils/errors'
import { isRecordLike } from '@sim/utils/object'
import { truncate } from '@sim/utils/string'
import { eq } from 'drizzle-orm'
import { type NextRequest, NextResponse } from 'next/server'
import { adminMothershipQuerySchema } from '@/lib/api/contracts/mothership-chats'
Expand All @@ -11,6 +14,10 @@ import { env } from '@/lib/core/config/env'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url'

const logger = createLogger('AdminMothershipProxy')

const UPSTREAM_BODY_LOG_LIMIT = 500

const ENV_URLS: Record<string, string | undefined> = {
dev: env.MOTHERSHIP_DEV_URL,
staging: env.MOTHERSHIP_STAGING_URL,
Expand Down Expand Up @@ -54,77 +61,106 @@ async function getAuthorizedAdminUserId() {
return authorized ? session.user.id : null
}

/**
* Proxy to the mothership admin API.
*
* Query params:
* env - "dev" | "staging" | "prod"
* endpoint - the admin endpoint path, e.g. "requests", "licenses", "traces"
*
* The request body (for POST) is forwarded as-is. Additional query params
* (e.g. requestId for GET /traces) are forwarded.
*/
export const POST = withRouteHandler(async (req: NextRequest) => {
const userId = await getAuthorizedAdminUserId()
if (!userId) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}

const adminKey = env.MOTHERSHIP_API_ADMIN_KEY
if (!adminKey) {
return NextResponse.json({ error: 'MOTHERSHIP_API_ADMIN_KEY not configured' }, { status: 500 })
}

const { searchParams } = new URL(req.url)
const queryValidation = adminMothershipQuerySchema.safeParse(searchParamsToObject(searchParams))
if (!queryValidation.success) return validationErrorResponse(queryValidation.error)
const { env: environment, endpoint } = queryValidation.data

if (!isValidEndpoint(endpoint)) {
return NextResponse.json({ error: 'invalid endpoint' }, { status: 400 })
/** `undefined` when the body is not JSON. Only called on a non-empty body. */
function parseJsonText(text: string): unknown {
try {
return JSON.parse(text)
} catch {
return undefined
}
}

const baseUrl = await getMothershipUrl(environment, userId)
if (!baseUrl) {
return NextResponse.json(
{ error: `No URL configured for environment: ${environment}` },
{ status: 400 }
)
}
/** The upstream's own explanation, from either error envelope the admin API uses. */
function upstreamErrorMessage(data: unknown): string | undefined {
if (!isRecordLike(data)) return undefined
if (typeof data.error === 'string') return data.error
if (typeof data.message === 'string') return data.message
return undefined
}

const targetUrl = `${baseUrl}/api/admin/${endpoint}`
type ProxyMethod = 'GET' | 'POST' | 'DELETE'

/**
* Forwards to the mothership admin API. A 4xx passes through for the admin UI to show; an
* upstream 5xx or unparseable body is a gateway failure, logged with its cause and answered
* with 502 so it is not mistaken for a failure of this route.
*/
async function forwardToMothership(params: {
method: ProxyMethod
targetUrl: string
adminKey: string
environment: string
endpoint: string
body?: string
}) {
const { method, targetUrl, adminKey, environment, endpoint, body } = params
try {
const body = await req.text()
const upstream = await fetch(targetUrl, {
method: 'POST',
method,
headers: {
'Content-Type': 'application/json',
...(method === 'POST' ? { 'Content-Type': 'application/json' } : {}),
'x-api-key': adminKey,
},
...(body ? { body } : {}),
})
const text = await upstream.text()
if (!text && upstream.status < 500) {
return new NextResponse(null, { status: upstream.status })
}
const data = text ? parseJsonText(text) : null

if (upstream.status >= 500 || data === undefined) {
logger.error('Mothership admin API request failed', {
method,
environment,
endpoint,
status: upstream.status,
body: truncate(text, UPSTREAM_BODY_LOG_LIMIT),
})
return NextResponse.json(
{
error:
upstreamErrorMessage(data) ??
`Mothership (${environment}) returned HTTP ${upstream.status}${data === undefined ? ' with a non-JSON body' : ''}`,
},
{ status: 502 }
)
}

const data = await upstream.json()
return NextResponse.json(data, { status: upstream.status })
} catch (error) {
logger.error('Failed to reach mothership admin API', {
method,
environment,
endpoint,
error: getErrorMessage(error, 'Unknown error'),
})
return NextResponse.json(
{
error: `Failed to reach mothership (${environment}): ${getErrorMessage(error, 'Unknown error')}`,
},
{ status: 502 }
)
}
})
}

export const GET = withRouteHandler(async (req: NextRequest) => {
/**
* Query params:
* env - "dev" | "staging" | "prod"
* endpoint - the admin endpoint path, e.g. "requests", "licenses", "traces"
*
* The request body (for POST) is forwarded as-is. For GET and DELETE, additional query
* params (e.g. requestId for GET /traces) are forwarded.
*/
async function proxyAdminRequest(req: NextRequest, method: ProxyMethod) {
const userId = await getAuthorizedAdminUserId()
if (!userId) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}

const adminKey = env.MOTHERSHIP_API_ADMIN_KEY
if (!adminKey) {
logger.error('MOTHERSHIP_API_ADMIN_KEY is not configured', { method })
return NextResponse.json({ error: 'MOTHERSHIP_API_ADMIN_KEY not configured' }, { status: 500 })
}

Expand All @@ -146,85 +182,25 @@ export const GET = withRouteHandler(async (req: NextRequest) => {
}

const forwardParams = new URLSearchParams()
searchParams.forEach((value, key) => {
if (key !== 'env' && key !== 'endpoint') {
forwardParams.set(key, value)
}
})

const qs = forwardParams.toString()
const targetUrl = `${baseUrl}/api/admin/${endpoint}${qs ? `?${qs}` : ''}`

try {
const upstream = await fetch(targetUrl, {
method: 'GET',
headers: { 'x-api-key': adminKey },
if (method !== 'POST') {
searchParams.forEach((value, key) => {
if (key !== 'env' && key !== 'endpoint') {
forwardParams.set(key, value)
}
})

const data = await upstream.json()
return NextResponse.json(data, { status: upstream.status })
} catch (error) {
return NextResponse.json(
{
error: `Failed to reach mothership (${environment}): ${getErrorMessage(error, 'Unknown error')}`,
},
{ status: 502 }
)
}
})

export const DELETE = withRouteHandler(async (req: NextRequest) => {
const userId = await getAuthorizedAdminUserId()
if (!userId) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}

const adminKey = env.MOTHERSHIP_API_ADMIN_KEY
if (!adminKey) {
return NextResponse.json({ error: 'MOTHERSHIP_API_ADMIN_KEY not configured' }, { status: 500 })
}

const { searchParams } = new URL(req.url)
const queryValidation = adminMothershipQuerySchema.safeParse(searchParamsToObject(searchParams))
if (!queryValidation.success) return validationErrorResponse(queryValidation.error)
const { env: environment, endpoint } = queryValidation.data

if (!isValidEndpoint(endpoint)) {
return NextResponse.json({ error: 'invalid endpoint' }, { status: 400 })
}

const baseUrl = await getMothershipUrl(environment, userId)
if (!baseUrl) {
return NextResponse.json(
{ error: `No URL configured for environment: ${environment}` },
{ status: 400 }
)
}

const forwardParams = new URLSearchParams()
searchParams.forEach((value, key) => {
if (key !== 'env' && key !== 'endpoint') {
forwardParams.set(key, value)
}
})

const qs = forwardParams.toString()
const targetUrl = `${baseUrl}/api/admin/${endpoint}${qs ? `?${qs}` : ''}`

try {
const upstream = await fetch(targetUrl, {
method: 'DELETE',
headers: { 'x-api-key': adminKey },
})
return forwardToMothership({
method,
targetUrl: `${baseUrl}/api/admin/${endpoint}${qs ? `?${qs}` : ''}`,
adminKey,
environment,
endpoint,
...(method === 'POST' ? { body: await req.text() } : {}),
})
}

const data = await upstream.json()
return NextResponse.json(data, { status: upstream.status })
} catch (error) {
return NextResponse.json(
{
error: `Failed to reach mothership (${environment}): ${getErrorMessage(error, 'Unknown error')}`,
},
{ status: 502 }
)
}
})
export const POST = withRouteHandler((req: NextRequest) => proxyAdminRequest(req, 'POST'))
export const GET = withRouteHandler((req: NextRequest) => proxyAdminRequest(req, 'GET'))
export const DELETE = withRouteHandler((req: NextRequest) => proxyAdminRequest(req, 'DELETE'))
2 changes: 1 addition & 1 deletion apps/sim/app/api/guardrails/mask-batch/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
try {
const startedAt = performance.now()
const masked = await maskPIIBatch(texts, entityTypes, language, customPatterns)
logger.info('Masked PII batch', {
logger.debug('Masked PII batch', {
count: texts.length,
durationMs: Math.round(performance.now() - startedAt),
})
Expand Down
9 changes: 6 additions & 3 deletions apps/sim/app/api/telemetry/route.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { createLogger } from '@sim/logger'
import { getErrorMessage } from '@sim/utils/errors'
import { type NextRequest, NextResponse } from 'next/server'
import { telemetryContract } from '@/lib/api/contracts/telemetry'
import { parseRequest } from '@/lib/api/server'
Expand Down Expand Up @@ -122,7 +123,7 @@ async function forwardToCollector(data: Record<string, unknown>): Promise<boolea
clearTimeout(timeoutId)

if (!response.ok) {
logger.error('Telemetry collector returned error', {
logger.warn('Telemetry collector returned error', {
status: response.status,
statusText: response.statusText,
})
Expand All @@ -133,9 +134,11 @@ async function forwardToCollector(data: Record<string, unknown>): Promise<boolea
} catch (fetchError) {
clearTimeout(timeoutId)
if (fetchError instanceof Error && fetchError.name === 'AbortError') {
logger.error('Telemetry request timed out', { endpoint })
logger.warn('Telemetry request timed out', { endpoint })
} else {
logger.error('Failed to send telemetry to collector', fetchError)
logger.warn('Failed to send telemetry to collector', {
error: getErrorMessage(fetchError),
})
}
return false
}
Expand Down
28 changes: 22 additions & 6 deletions apps/sim/background/webhook-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -722,17 +722,32 @@ export async function resolveWebhookExecutionProviderConfig<
}
}

async function resolveCredentialAccountUserId(credentialId: string): Promise<string | undefined> {
/**
* The user who owns the OAuth account behind a credential. `accountless` credentials
* (service accounts, managed OAuth such as a custom Slack bot) have no account row by
* design; `missing` means the credential or its account no longer exists.
*/
async function resolveCredentialAccountUserId(
credentialId: string
): Promise<{ status: 'resolved'; userId: string } | { status: 'accountless' | 'missing' }> {
const resolved = await resolveOAuthAccountId(credentialId)
if (!resolved) {
return undefined
return { status: 'missing' }
}
if (
resolved.credentialType === 'service_account' ||
resolved.credentialType === 'managed_oauth'
) {
return { status: 'accountless' }
}
const [credentialRecord] = await db
.select({ userId: account.userId })
.from(account)
.where(eq(account.id, resolved.accountId))
.limit(1)
return credentialRecord?.userId
return credentialRecord
? { status: 'resolved', userId: credentialRecord.userId }
: { status: 'missing' }
}

/**
Expand Down Expand Up @@ -873,15 +888,16 @@ async function executeWebhookJobInternal(
workspaceId
)
: loadDeployedWorkflowState(payload.workflowId, workspaceId)
const [workflowData, webhookRows, resolvedCredentialUserId] = await Promise.all([
const [workflowData, webhookRows, credentialAccount] = await Promise.all([
workflowStatePromise,
db.select().from(webhook).where(eq(webhook.id, payload.webhookId)).limit(1),
payload.credentialId
? resolveCredentialAccountUserId(payload.credentialId)
: Promise.resolve(undefined),
])
const credentialAccountUserId = resolvedCredentialUserId
if (payload.credentialId && !credentialAccountUserId) {
const credentialAccountUserId =
credentialAccount?.status === 'resolved' ? credentialAccount.userId : undefined
if (credentialAccount?.status === 'missing') {
logger.warn(
`[${requestId}] Failed to resolve credential account for credential ${payload.credentialId}`
)
Expand Down
2 changes: 1 addition & 1 deletion apps/sim/blocks/custom/server-overlay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { registerBlockOverlayResolver } from '@/blocks/custom/overlay'
import type { BlockConfig, BlockIcon } from '@/blocks/types'

/** A row for the overlay, optionally carrying live-derived Start input fields. */
type CustomBlockOverlayRow = CustomBlockRow & {
export type CustomBlockOverlayRow = CustomBlockRow & {
inputFields?: WorkflowInputField[]
/** When `false`, the block resolves but is hidden from the palette (disabled). */
enabled?: boolean
Expand Down
2 changes: 1 addition & 1 deletion apps/sim/executor/dag/builder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ export class DAGBuilder {
// Validate loop and parallel structure
this.validateSubflowStructure(dag)

logger.info('DAG built', {
logger.debug('DAG built', {
totalNodes: dag.nodes.size,
loopCount: dag.loopConfigs.size,
parallelCount: dag.parallelConfigs.size,
Expand Down
Loading
Loading