Files
WeHub Mirror 6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

469 lines
15 KiB
TypeScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import type { Logger } from '@sim/logger'
import { getErrorMessage } from '@sim/utils/errors'
import { pollingIdempotency } from '@/lib/core/idempotency/service'
import { readCanonicalTriggerValue } from '@/lib/webhooks/polling/canonical'
import {
getProviderConfig,
type PollingProviderHandler,
type PollWebhookContext,
} from '@/lib/webhooks/polling/types'
import {
markWebhookFailed,
markWebhookSuccess,
resolveOAuthCredential,
updateWebhookProviderConfig,
} from '@/lib/webhooks/polling/utils'
import { processPolledWebhookEvent } from '@/lib/webhooks/processor'
const MAX_ROWS_PER_POLL = 100
/** Maximum number of leading rows to scan when auto-detecting the header row. */
const HEADER_SCAN_ROWS = 10
type ValueRenderOption = 'FORMATTED_VALUE' | 'UNFORMATTED_VALUE' | 'FORMULA'
type DateTimeRenderOption = 'SERIAL_NUMBER' | 'FORMATTED_STRING'
interface GoogleSheetsWebhookConfig {
spreadsheetId?: string
manualSpreadsheetId?: string
sheetName?: string
manualSheetName?: string
valueRenderOption?: ValueRenderOption
dateTimeRenderOption?: DateTimeRenderOption
/** 1-indexed row number of the last row seeded or processed. */
lastIndexChecked?: number
lastModifiedTime?: string
lastCheckedTimestamp?: string
maxRowsPerPoll?: number
}
interface GoogleSheetsWebhookPayload {
row: Record<string, string> | null
rawRow: string[]
headers: string[]
rowNumber: number
spreadsheetId: string
sheetName: string
timestamp: string
}
export const googleSheetsPollingHandler: PollingProviderHandler = {
provider: 'google-sheets',
label: 'Google Sheets',
async pollWebhook(ctx: PollWebhookContext): Promise<'success' | 'failure'> {
const { webhookData, workflowData, requestId, logger } = ctx
const webhookId = webhookData.id
try {
const accessToken = await resolveOAuthCredential(webhookData, 'google-sheets', requestId)
const config = getProviderConfig<GoogleSheetsWebhookConfig>(webhookData.providerConfig)
// Canonical keys (`spreadsheetId`/`sheetName`) first; the `manual*` keys are a transitional
// basic-first fallback for configs deployed before the canonical key existed.
const spreadsheetId = readCanonicalTriggerValue(
config.spreadsheetId,
config.manualSpreadsheetId
)
const sheetName = readCanonicalTriggerValue(config.sheetName, config.manualSheetName)
const now = new Date()
if (!spreadsheetId || !sheetName) {
logger.error(`[${requestId}] Missing spreadsheetId or sheetName for webhook ${webhookId}`)
await markWebhookFailed(webhookId, logger)
return 'failure'
}
const { unchanged: skipPoll, currentModifiedTime } = await isDriveFileUnchanged(
accessToken,
spreadsheetId,
config.lastModifiedTime,
requestId,
logger
)
if (skipPoll) {
await updateWebhookProviderConfig(
webhookId,
{ lastCheckedTimestamp: now.toISOString() },
logger
)
await markWebhookSuccess(webhookId, logger)
logger.info(`[${requestId}] Sheet not modified since last poll for webhook ${webhookId}`)
return 'success'
}
const valueRender = config.valueRenderOption || 'FORMATTED_VALUE'
const dateTimeRender = config.dateTimeRenderOption || 'SERIAL_NUMBER'
const {
rowCount: currentRowCount,
headers,
headerRowIndex,
} = await fetchSheetState(
accessToken,
spreadsheetId,
sheetName,
valueRender,
dateTimeRender,
requestId,
logger
)
// First poll: seed state, emit nothing
if (config.lastIndexChecked === undefined) {
await updateWebhookProviderConfig(
webhookId,
{
lastIndexChecked: currentRowCount,
lastModifiedTime: currentModifiedTime ?? config.lastModifiedTime,
lastCheckedTimestamp: now.toISOString(),
},
logger
)
await markWebhookSuccess(webhookId, logger)
logger.info(
`[${requestId}] First poll for webhook ${webhookId}, seeded row index: ${currentRowCount}`
)
return 'success'
}
if (currentRowCount <= config.lastIndexChecked) {
if (currentRowCount < config.lastIndexChecked) {
logger.warn(
`[${requestId}] Row count decreased from ${config.lastIndexChecked} to ${currentRowCount} for webhook ${webhookId}`
)
}
await updateWebhookProviderConfig(
webhookId,
{
lastIndexChecked: currentRowCount,
lastModifiedTime: currentModifiedTime ?? config.lastModifiedTime,
lastCheckedTimestamp: now.toISOString(),
},
logger
)
await markWebhookSuccess(webhookId, logger)
logger.info(`[${requestId}] No new rows for webhook ${webhookId}`)
return 'success'
}
const newRowCount = currentRowCount - config.lastIndexChecked
const maxRows = config.maxRowsPerPoll || MAX_ROWS_PER_POLL
const rowsToFetch = Math.min(newRowCount, maxRows)
const startRow = config.lastIndexChecked + 1
const endRow = config.lastIndexChecked + rowsToFetch
// Skip past the header row (and any blank rows above it) so it is never
// emitted as a data event.
const adjustedStartRow =
headerRowIndex > 0 ? Math.max(startRow, headerRowIndex + 1) : startRow
logger.info(
`[${requestId}] Found ${newRowCount} new rows for webhook ${webhookId}, processing rows ${adjustedStartRow}-${endRow}`
)
// Entire batch is header/blank rows — advance pointer and skip fetch.
if (adjustedStartRow > endRow) {
const hasRemainingRows = rowsToFetch < newRowCount
await updateWebhookProviderConfig(
webhookId,
{
lastIndexChecked: config.lastIndexChecked + rowsToFetch,
lastModifiedTime: hasRemainingRows
? config.lastModifiedTime
: (currentModifiedTime ?? config.lastModifiedTime),
lastCheckedTimestamp: now.toISOString(),
},
logger
)
await markWebhookSuccess(webhookId, logger)
logger.info(
`[${requestId}] Batch ${startRow}-${endRow} contained only header/blank rows for webhook ${webhookId}, advancing pointer`
)
return 'success'
}
const newRows = await fetchRowRange(
accessToken,
spreadsheetId,
sheetName,
adjustedStartRow,
endRow,
valueRender,
dateTimeRender,
requestId,
logger
)
const { processedCount, failedCount } = await processRows(
newRows,
headers,
adjustedStartRow,
spreadsheetId,
sheetName,
webhookData,
workflowData,
requestId,
logger
)
const rowsAdvanced = failedCount > 0 ? 0 : rowsToFetch
const newLastIndexChecked = config.lastIndexChecked + rowsAdvanced
const hasRemainingOrFailed = rowsAdvanced < newRowCount
await updateWebhookProviderConfig(
webhookId,
{
lastIndexChecked: newLastIndexChecked,
lastModifiedTime: hasRemainingOrFailed
? config.lastModifiedTime
: (currentModifiedTime ?? config.lastModifiedTime),
lastCheckedTimestamp: now.toISOString(),
},
logger
)
if (failedCount > 0 && processedCount === 0) {
await markWebhookFailed(webhookId, logger)
logger.warn(
`[${requestId}] All ${failedCount} rows failed to process for webhook ${webhookId}`
)
return 'failure'
}
await markWebhookSuccess(webhookId, logger)
logger.info(
`[${requestId}] Successfully processed ${processedCount} rows for webhook ${webhookId}${failedCount > 0 ? ` (${failedCount} failed)` : ''}`
)
return 'success'
} catch (error) {
logger.error(`[${requestId}] Error processing Google Sheets webhook ${webhookId}:`, error)
await markWebhookFailed(webhookId, logger)
return 'failure'
}
},
}
async function isDriveFileUnchanged(
accessToken: string,
spreadsheetId: string,
lastModifiedTime: string | undefined,
requestId: string,
logger: Logger
): Promise<{ unchanged: boolean; currentModifiedTime?: string }> {
try {
const currentModifiedTime = await getDriveFileModifiedTime(accessToken, spreadsheetId, logger)
if (!lastModifiedTime || !currentModifiedTime) {
return { unchanged: false, currentModifiedTime }
}
return { unchanged: currentModifiedTime === lastModifiedTime, currentModifiedTime }
} catch (error) {
logger.warn(`[${requestId}] Drive modifiedTime check failed, proceeding with Sheets API`)
return { unchanged: false }
}
}
async function getDriveFileModifiedTime(
accessToken: string,
fileId: string,
logger: Logger
): Promise<string | undefined> {
try {
const response = await fetch(
`https://www.googleapis.com/drive/v3/files/${fileId}?fields=modifiedTime`,
{ headers: { Authorization: `Bearer ${accessToken}` } }
)
if (!response.ok) return undefined
const data = await response.json()
return data.modifiedTime as string | undefined
} catch {
return undefined
}
}
/**
* Fetches the sheet (A:Z) and returns the row count, auto-detected headers,
* and the 1-indexed header row number in a single API call.
*
* The Sheets API omits trailing empty rows, so `rows.length` equals the last
* non-empty row in columns AZ. Header detection scans the first
* {@link HEADER_SCAN_ROWS} rows for the first non-empty row. Returns
* `headerRowIndex = 0` when no header is found within the scan window.
*/
async function fetchSheetState(
accessToken: string,
spreadsheetId: string,
sheetName: string,
valueRenderOption: ValueRenderOption,
dateTimeRenderOption: DateTimeRenderOption,
requestId: string,
logger: Logger
): Promise<{ rowCount: number; headers: string[]; headerRowIndex: number }> {
const encodedSheet = encodeURIComponent(sheetName)
const params = new URLSearchParams({
majorDimension: 'ROWS',
fields: 'values',
valueRenderOption,
dateTimeRenderOption,
})
const url = `https://sheets.googleapis.com/v4/spreadsheets/${spreadsheetId}/values/${encodedSheet}!A:Z?${params.toString()}`
const response = await fetch(url, {
headers: { Authorization: `Bearer ${accessToken}` },
})
if (!response.ok) {
const status = response.status
const errorData = await response.json().catch(() => ({}))
if (status === 403 || status === 429) {
throw new Error(
`Sheets API rate limit (${status}) — skipping to retry next poll cycle: ${JSON.stringify(errorData)}`
)
}
throw new Error(
`Failed to fetch sheet state: ${status} ${response.statusText} - ${JSON.stringify(errorData)}`
)
}
const data = await response.json()
const rows = (data.values as string[][] | undefined) ?? []
const rowCount = rows.length
let headers: string[] = []
let headerRowIndex = 0
for (let i = 0; i < Math.min(rows.length, HEADER_SCAN_ROWS); i++) {
const row = rows[i]
if (row?.some((cell) => cell !== '')) {
headers = row
headerRowIndex = i + 1
break
}
}
return { rowCount, headers, headerRowIndex }
}
async function fetchRowRange(
accessToken: string,
spreadsheetId: string,
sheetName: string,
startRow: number,
endRow: number,
valueRenderOption: ValueRenderOption,
dateTimeRenderOption: DateTimeRenderOption,
requestId: string,
logger: Logger
): Promise<string[][]> {
const encodedSheet = encodeURIComponent(sheetName)
const params = new URLSearchParams({
fields: 'values',
valueRenderOption,
dateTimeRenderOption,
})
const url = `https://sheets.googleapis.com/v4/spreadsheets/${spreadsheetId}/values/${encodedSheet}!${startRow}:${endRow}?${params.toString()}`
const response = await fetch(url, {
headers: { Authorization: `Bearer ${accessToken}` },
})
if (!response.ok) {
const status = response.status
const errorData = await response.json().catch(() => ({}))
if (status === 403 || status === 429) {
throw new Error(
`Sheets API rate limit (${status}) — skipping to retry next poll cycle: ${JSON.stringify(errorData)}`
)
}
throw new Error(
`Failed to fetch rows ${startRow}-${endRow}: ${status} ${response.statusText} - ${JSON.stringify(errorData)}`
)
}
const data = await response.json()
return (data.values as string[][]) ?? []
}
async function processRows(
rows: string[][],
headers: string[],
startRowIndex: number,
spreadsheetId: string,
sheetName: string,
webhookData: PollWebhookContext['webhookData'],
workflowData: PollWebhookContext['workflowData'],
requestId: string,
logger: Logger
): Promise<{ processedCount: number; failedCount: number }> {
let processedCount = 0
let failedCount = 0
for (let i = 0; i < rows.length; i++) {
const row = rows[i]
const rowNumber = startRowIndex + i
// Skip empty rows — don't fire a workflow run with no data.
if (!row || row.length === 0) {
logger.info(`[${requestId}] Skipping empty row ${rowNumber} for webhook ${webhookData.id}`)
continue
}
try {
await pollingIdempotency.executeWithIdempotency(
'google-sheets',
`${webhookData.id}:${spreadsheetId}:${sheetName}:row${rowNumber}`,
async () => {
let mappedRow: Record<string, string> | null = null
if (headers.length > 0) {
mappedRow = {}
for (let j = 0; j < headers.length; j++) {
mappedRow[headers[j] || `Column ${j + 1}`] = row[j] ?? ''
}
for (let j = headers.length; j < row.length; j++) {
mappedRow[`Column ${j + 1}`] = row[j] ?? ''
}
}
const payload: GoogleSheetsWebhookPayload = {
row: mappedRow,
rawRow: row,
headers,
rowNumber,
spreadsheetId,
sheetName,
timestamp: new Date().toISOString(),
}
const result = await processPolledWebhookEvent(
webhookData,
workflowData,
payload,
requestId
)
if (!result.success) {
logger.error(
`[${requestId}] Failed to process webhook for row ${rowNumber}:`,
result.statusCode,
result.error
)
throw new Error(`Webhook processing failed (${result.statusCode}): ${result.error}`)
}
return { rowNumber, processed: true }
}
)
logger.info(
`[${requestId}] Successfully processed row ${rowNumber} for webhook ${webhookData.id}`
)
processedCount++
} catch (error) {
const errorMessage = getErrorMessage(error, 'Unknown error')
logger.error(`[${requestId}] Error processing row ${rowNumber}:`, errorMessage)
failedCount++
}
}
return { processedCount, failedCount }
}