diff --git a/apps/sim/app/workspace/[workspaceId]/components/resource/components/resource-header/resource-header.tsx b/apps/sim/app/workspace/[workspaceId]/components/resource/components/resource-header/resource-header.tsx
index 006271eb23a..301cc2e7dba 100644
--- a/apps/sim/app/workspace/[workspaceId]/components/resource/components/resource-header/resource-header.tsx
+++ b/apps/sim/app/workspace/[workspaceId]/components/resource/components/resource-header/resource-header.tsx
@@ -39,7 +39,7 @@ import {
import { ArrowUpLeft } from '@sim/emcn/icons'
import { createPortal } from 'react-dom'
import { HEADER_ACTION_CLUSTER, TITLE_BAR_LANE_PT } from '@/components/page-header-bar'
-import { orderHeaderActions } from '@/components/settings/settings-header'
+import { orderHeaderActions, SettingsActionChip } from '@/components/settings/settings-header'
import { InlineRenameInput } from '@/app/workspace/[workspaceId]/components/inline-rename-input'
export interface DropdownOption {
@@ -104,6 +104,7 @@ export interface ResourceAction {
active?: boolean
onSelect: () => void
disabled?: boolean
+ tooltip?: string
}
/**
@@ -259,16 +260,7 @@ export const ResourceHeader = memo(function ResourceHeader({
{aside}
{orderHeaderActions(actions).map(({ action }) => (
-
- {action.text}
-
+
))}
)}
diff --git a/apps/sim/app/workspace/[workspaceId]/logs/components/log-refresh-action.tsx b/apps/sim/app/workspace/[workspaceId]/logs/components/log-refresh-action.tsx
new file mode 100644
index 00000000000..d8c89fe59e1
--- /dev/null
+++ b/apps/sim/app/workspace/[workspaceId]/logs/components/log-refresh-action.tsx
@@ -0,0 +1,55 @@
+import { cn } from '@sim/emcn'
+import { RefreshCw } from '@sim/emcn/icons'
+import type { ResourceAction } from '@/app/workspace/[workspaceId]/components/resource/components/resource-header'
+
+interface RefreshIconProps {
+ className?: string
+}
+
+function SpinningRefreshIcon(props: RefreshIconProps) {
+ return
+}
+
+function NewLogsIndicator({ className }: RefreshIconProps) {
+ return (
+
+
+
+ )
+}
+
+interface LogsRefreshActionOptions {
+ newLogCount: number
+ hasUpdates?: boolean
+ isRefreshing: boolean
+ onRefresh: () => void
+}
+
+export function getLogsRefreshAction({
+ newLogCount,
+ hasUpdates = false,
+ isRefreshing,
+ onRefresh,
+}: LogsRefreshActionOptions): ResourceAction {
+ const hasNewLogs = newLogCount > 0
+ return {
+ id: 'refresh',
+ text: hasNewLogs
+ ? `${newLogCount} new ${newLogCount === 1 ? 'log' : 'logs'}`
+ : hasUpdates
+ ? 'Updates available'
+ : 'Refresh',
+ tooltip: hasNewLogs
+ ? 'Refresh to see new logs'
+ : hasUpdates
+ ? 'Refresh to see updated logs'
+ : 'Refresh',
+ icon: isRefreshing
+ ? SpinningRefreshIcon
+ : hasNewLogs || hasUpdates
+ ? NewLogsIndicator
+ : RefreshCw,
+ onSelect: onRefresh,
+ disabled: isRefreshing,
+ }
+}
diff --git a/apps/sim/app/workspace/[workspaceId]/logs/logs.tsx b/apps/sim/app/workspace/[workspaceId]/logs/logs.tsx
index 7974728b3e7..b93297fe27a 100644
--- a/apps/sim/app/workspace/[workspaceId]/logs/logs.tsx
+++ b/apps/sim/app/workspace/[workspaceId]/logs/logs.tsx
@@ -22,7 +22,6 @@ import {
Popover,
PopoverAnchor,
PopoverContent,
- RefreshCw,
toast,
} from '@sim/emcn'
import { Download, Workflow } from '@sim/emcn/icons'
@@ -73,6 +72,7 @@ import {
SnapshotBoundary,
SnapshotModalFallback,
} from '@/app/workspace/[workspaceId]/logs/components/log-details/components/execution-snapshot/snapshot-boundary'
+import { getLogsRefreshAction } from '@/app/workspace/[workspaceId]/logs/components/log-refresh-action'
import { useLogFilters } from '@/app/workspace/[workspaceId]/logs/hooks/use-log-filters'
import { useSearchState } from '@/app/workspace/[workspaceId]/logs/hooks/use-search-state'
import {
@@ -94,7 +94,9 @@ import {
useDashboardStats,
useLogByExecutionId,
useLogDetail,
- useLogsList,
+ useLogSnapshotUpdates,
+ useLogsSnapshot,
+ useNewLogCount,
useRetryExecution,
} from '@/hooks/queries/logs'
import { useWorkflowMap, useWorkflows } from '@/hooks/queries/workflows'
@@ -130,7 +132,6 @@ const ExecutionSnapshot = lazy(() =>
)
const LOGS_PER_PAGE = 50 as const
-const REFRESH_SPINNER_DURATION_MS = 1000 as const
const LIVE_REFRESH_INTERVAL_MS = 10_000 as const
const ACTIVE_RUN_DETAIL_REFRESH_MS = 3_000 as const
@@ -228,8 +229,11 @@ function getTriggerIcon(
return TriggerIcon
}
-function SpinningRefreshCw(props: React.SVGProps) {
- return
+function activeRunRefetchInterval(query: { state: { data?: WorkflowLogDetail } }) {
+ const status = query.state.data?.status
+ return status === 'running' || status === 'pending' || status === 'redacting'
+ ? ACTIVE_RUN_DETAIL_REFRESH_MS
+ : false
}
/**
@@ -291,17 +295,14 @@ export default function Logs() {
*/
const debouncedSearchQuery = useDebounce(urlSearchQuery, SEARCH_DEBOUNCE_MS).trim()
- const isLive = true
- const [isVisuallyRefreshing, setIsVisuallyRefreshing] = useState(false)
const [isExporting, setIsExporting] = useState(false)
- const refreshTimersRef = useRef(new Set())
const logsRef = useRef([])
const selectedLogIndexRef = useRef(-1)
const selectedLogIdRef = useRef(null)
const isSidebarOpenRef = useRef(false)
const shouldScrollIntoViewRef = useRef(false)
const resourceTableRef = useRef(null)
- const activeViewRefetchRef = useRef<() => void>(() => {})
+ const activeViewRefetchRef = useRef<() => Promise>(async () => {})
const activeLogRefetchRef = useRef<() => void>(() => {})
const activeLogTabRef = useRef('overview')
const logsQueryRef = useRef({ isFetching: false, hasNextPage: false, fetchNextPage: () => {} })
@@ -328,24 +329,13 @@ export default function Logs() {
const queryClient = useQueryClient()
- const refetchInterval = useCallback(
- (query: { state: { data?: WorkflowLogDetail } }) => {
- if (!isLive) return false
- const status = query.state.data?.status
- return status === 'running' || status === 'pending' || status === 'redacting'
- ? ACTIVE_RUN_DETAIL_REFRESH_MS
- : false
- },
- [isLive]
- )
-
const selectedDetailQuery = useLogDetail(selectedLogId ?? undefined, workspaceId, {
enabled: isSidebarOpen,
- refetchInterval,
+ refetchInterval: activeRunRefetchInterval,
})
const previewDetailQuery = useLogDetail(previewLogId ?? undefined, workspaceId, {
- refetchInterval,
+ refetchInterval: activeRunRefetchInterval,
})
const logFilters = useMemo(
@@ -376,10 +366,23 @@ export default function Logs() {
]
)
- const logsQuery = useLogsList(workspaceId, logFilters, {
+ const logsQuery = useLogsSnapshot(workspaceId, logFilters, {
enabled: !isDashboardView || isSidebarOpen,
- refetchInterval: isLive ? LIVE_REFRESH_INTERVAL_MS : false,
})
+ const newLogsQuery = useNewLogCount(
+ workspaceId,
+ logFilters,
+ logsQuery.isPlaceholderData ? undefined : logsQuery.data?.pages[0]?.snapshotAt,
+ { enabled: !isDashboardView }
+ )
+ const newLogCount = newLogsQuery.data ?? 0
+ const hasChangedPage = logsQuery.data?.pages.some((page) => page.snapshotChanged) === true
+ const snapshotUpdatesQuery = useLogSnapshotUpdates(
+ workspaceId,
+ logsQuery.isPlaceholderData ? undefined : logsQuery.data?.pages[0],
+ { enabled: !isDashboardView && newLogCount === 0 && !hasChangedPage }
+ )
+ const hasSnapshotUpdates = snapshotUpdatesQuery.data || hasChangedPage
const dashboardFilters = useMemo(
() => ({
@@ -397,7 +400,7 @@ export default function Logs() {
const dashboardStatsQuery = useDashboardStats(workspaceId, dashboardFilters, {
enabled: isDashboardView,
- refetchInterval: isLive ? LIVE_REFRESH_INTERVAL_MS : false,
+ refetchInterval: LIVE_REFRESH_INTERVAL_MS,
})
const logs = useMemo(() => {
@@ -421,14 +424,11 @@ export default function Logs() {
selectedLogIndexRef.current = selectedLogIndex
selectedLogIdRef.current = selectedLogId
isSidebarOpenRef.current = isSidebarOpen
- activeViewRefetchRef.current = () => {
- if (isDashboardView) {
- void dashboardStatsQuery.refetch()
- }
- if (!isDashboardView || isSidebarOpen) {
- void logsQuery.refetch()
- }
- }
+ activeViewRefetchRef.current = () =>
+ Promise.all([
+ ...(isDashboardView ? [dashboardStatsQuery.refetch({ throwOnError: true })] : []),
+ ...(!isDashboardView || isSidebarOpen ? [logsQuery.refetch({ throwOnError: true })] : []),
+ ])
activeLogRefetchRef.current = selectedDetailQuery.refetch
logsQueryRef.current = {
isFetching: logsQuery.isFetching,
@@ -449,14 +449,6 @@ export default function Logs() {
}
}, [pendingExecutionId, deepLinkQuery.data, deepLinkQuery.isError])
- useEffect(() => {
- const timers = refreshTimersRef.current
- return () => {
- timers.forEach((id) => window.clearTimeout(id))
- timers.clear()
- }
- }, [])
-
/**
* The single write path for user-driven `executionId` changes. Cancels any
* in-flight deep-link resolution first — an explicit interaction supersedes
@@ -664,36 +656,20 @@ export default function Logs() {
const effectiveSidebarOpen =
isSidebarOpen && (selectedLogIndex !== -1 || !!selectedDetailQuery.data)
- const triggerVisualRefresh = useCallback(() => {
- setIsVisuallyRefreshing(true)
- const timerId = window.setTimeout(() => {
- setIsVisuallyRefreshing(false)
- refreshTimersRef.current.delete(timerId)
- }, REFRESH_SPINNER_DURATION_MS)
- refreshTimersRef.current.add(timerId)
- }, [])
-
- const handleRefresh = useCallback(() => {
- triggerVisualRefresh()
- activeViewRefetchRef.current()
+ const handleRefresh = useCallback(async () => {
if (selectedLogIdRef.current && isSidebarOpenRef.current) {
activeLogRefetchRef.current()
}
- }, [triggerVisualRefresh])
+ try {
+ await activeViewRefetchRef.current()
+ } catch (error) {
+ toast.error(getErrorMessage(error, 'Failed to refresh logs'))
+ }
+ }, [])
- const activeViewIsFetching = isDashboardView
+ const isVisuallyRefreshing = isDashboardView
? dashboardStatsQuery.isFetching || (isSidebarOpen && logsQuery.isFetching)
: logsQuery.isFetching
- const prevIsFetchingRef = useRef(activeViewIsFetching)
- useEffect(() => {
- const wasFetching = prevIsFetchingRef.current
- const isFetching = activeViewIsFetching
- prevIsFetchingRef.current = isFetching
-
- if (isLive && !wasFetching && isFetching) {
- triggerVisualRefresh()
- }
- }, [activeViewIsFetching, isLive, triggerVisualRefresh])
const handleExport = useCallback(async () => {
setIsExporting(true)
@@ -1169,7 +1145,6 @@ export default function Logs() {
]
)
- const refreshIcon = isVisuallyRefreshing ? SpinningRefreshCw : RefreshCw
const hasExportableLogs = isDashboardView
? !dashboardStatsQuery.isPlaceholderData && (dashboardStatsQuery.data?.totalRuns ?? 0) > 0
: !logsQuery.isPlaceholderData && logs.length > 0
@@ -1182,12 +1157,12 @@ export default function Logs() {
onSelect: handleExport,
disabled: !userPermissions.canEdit || isExporting || !hasExportableLogs,
},
- {
- text: 'Refresh',
- icon: refreshIcon,
- onSelect: handleRefresh,
- disabled: isVisuallyRefreshing,
- },
+ getLogsRefreshAction({
+ newLogCount: isDashboardView ? 0 : newLogCount,
+ hasUpdates: !isDashboardView && hasSnapshotUpdates,
+ isRefreshing: isVisuallyRefreshing,
+ onRefresh: handleRefresh,
+ }),
{
text: 'Logs',
onSelect: () => setViewMode('logs'),
@@ -1202,7 +1177,8 @@ export default function Logs() {
[
isDashboardView,
setViewMode,
- refreshIcon,
+ newLogCount,
+ hasSnapshotUpdates,
isVisuallyRefreshing,
handleRefresh,
handleExport,
diff --git a/apps/sim/hooks/queries/logs.test.tsx b/apps/sim/hooks/queries/logs.test.tsx
index 7ec36e6bf3b..bbb8ec6eacb 100644
--- a/apps/sim/hooks/queries/logs.test.tsx
+++ b/apps/sim/hooks/queries/logs.test.tsx
@@ -2,7 +2,12 @@
* @vitest-environment jsdom
*/
import { act, type ReactNode } from 'react'
-import { QueryClient, QueryClientProvider } from '@tanstack/react-query'
+import {
+ focusManager,
+ onlineManager,
+ QueryClient,
+ QueryClientProvider,
+} from '@tanstack/react-query'
import { createRoot, type Root } from 'react-dom/client'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
@@ -16,11 +21,22 @@ vi.mock('@/lib/api/client/request', () => ({
import { getLogByExecutionIdContract } from '@/lib/api/contracts/logs'
import { cancelWorkflowExecutionContract } from '@/lib/api/contracts/workflows'
-import { useCancelExecution } from '@/hooks/queries/logs'
+import {
+ LOG_SNAPSHOT_UPDATES_STALE_TIME,
+ type LogFilters,
+ logKeys,
+ NEW_LOG_COUNT_STALE_TIME,
+ useCancelExecution,
+ useLogSnapshotUpdates,
+ useLogsSnapshot,
+ useNewLogCount,
+} from '@/hooks/queries/logs'
function renderHookWithClient(useHook: () => T): {
result: () => T
unmount: () => void
+ rerender: () => void
+ queryClient: QueryClient
} {
;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true
const queryClient = new QueryClient({
@@ -52,10 +68,348 @@ function renderHookWithClient(useHook: () => T): {
return {
result: () => latest,
- unmount: () => act(() => root.unmount()),
+ unmount: () => {
+ act(() => root.unmount())
+ queryClient.clear()
+ },
+ rerender: () =>
+ act(() =>
+ root.render(
+
+
+
+ )
+ ),
+ queryClient,
}
}
+const SNAPSHOT_AT = '2026-09-24T15:45:00.000Z'
+const NEXT_SNAPSHOT_AT = '2026-09-24T15:46:00.000Z'
+const LOG_FILTERS: LogFilters = {
+ timeRange: 'All time',
+ level: 'all',
+ workflowIds: [],
+ folderIds: [],
+ triggers: [],
+ searchQuery: '',
+ limit: 50,
+ sortBy: 'date',
+ sortOrder: 'desc',
+}
+
+async function flushQueries() {
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(10)
+ })
+}
+
+describe('manually refreshed logs', () => {
+ let unmount: (() => void) | undefined
+ let snapshotAt: string
+ let newCount: number
+ let failRefresh: boolean
+ let revision: string
+
+ beforeEach(() => {
+ vi.clearAllMocks()
+ vi.useFakeTimers()
+ snapshotAt = SNAPSHOT_AT
+ newCount = 0
+ failRefresh = false
+ revision = 'first-revision'
+ mockRequestJson.mockImplementation(async (_contract, { query }) => {
+ if (query.startedAfter) return { data: [], nextCursor: null, total: newCount }
+ if (query.countOnly) return { data: [], nextCursor: null, total: 1, revision }
+ if (failRefresh) throw new Error('Refresh failed')
+ return {
+ data: [{ id: query.cursor ? 'older-log' : snapshotAt, status: 'running' }],
+ nextCursor: query.cursor ? null : 'next-page',
+ snapshotAt,
+ revision,
+ }
+ })
+ })
+
+ afterEach(() => {
+ unmount?.()
+ unmount = undefined
+ focusManager.setFocused(undefined)
+ onlineManager.setOnline(true)
+ vi.useRealTimers()
+ })
+
+ function useLogs(
+ filters = LOG_FILTERS,
+ enabled = true,
+ workspaceId: string | undefined = 'ws-1'
+ ) {
+ const list = useLogsSnapshot(workspaceId, filters, { enabled })
+ const count = useNewLogCount(
+ workspaceId,
+ filters,
+ list.isPlaceholderData ? undefined : list.data?.pages[0]?.snapshotAt,
+ { enabled }
+ )
+ const updates = useLogSnapshotUpdates(
+ workspaceId,
+ list.isPlaceholderData ? undefined : list.data?.pages[0],
+ { enabled: enabled && (count.data ?? 0) === 0 }
+ )
+ return { list, count, updates }
+ }
+
+ it('polls for new logs without changing rows on polling, focus, reconnect, or invalidation', async () => {
+ const hook = renderHookWithClient(() => useLogs())
+ unmount = hook.unmount
+ await flushQueries()
+ const originalRows = hook.result().list.data
+ newCount = 3
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(NEW_LOG_COUNT_STALE_TIME)
+ focusManager.setFocused(false)
+ focusManager.setFocused(true)
+ onlineManager.setOnline(false)
+ onlineManager.setOnline(true)
+ await hook.queryClient.invalidateQueries({ queryKey: logKeys.all })
+ })
+ await flushQueries()
+
+ expect(hook.result().count.data).toBe(3)
+ expect(hook.result().list.data).toBe(originalRows)
+ expect(mockRequestJson.mock.calls.filter(([, { query }]) => !query.startedAfter)).toHaveLength(
+ 1
+ )
+ })
+
+ it('pins pagination to the displayed snapshot and acknowledges new logs only after refresh', async () => {
+ const hook = renderHookWithClient(() => useLogs())
+ unmount = hook.unmount
+ await flushQueries()
+ newCount = 2
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(NEW_LOG_COUNT_STALE_TIME)
+ })
+ await act(async () => {
+ await hook.result().list.fetchNextPage()
+ })
+ await flushQueries()
+
+ expect(hook.result().list.data?.pages).toHaveLength(2)
+ expect(hook.result().count.data).toBe(2)
+ expect(mockRequestJson).toHaveBeenCalledWith(
+ expect.anything(),
+ expect.objectContaining({
+ query: expect.objectContaining({ cursor: 'next-page', snapshotAt: SNAPSHOT_AT }),
+ signal: expect.any(AbortSignal),
+ })
+ )
+
+ snapshotAt = NEXT_SNAPSHOT_AT
+ newCount = 0
+ await act(async () => {
+ await hook.result().list.refetch()
+ })
+ await flushQueries()
+
+ expect(hook.result().list.data?.pages[0]?.logs[0]?.id).toBe(NEXT_SNAPSHOT_AT)
+ expect(hook.result().count.data).toBe(0)
+ expect(mockRequestJson).toHaveBeenCalledWith(
+ expect.anything(),
+ expect.objectContaining({
+ query: expect.objectContaining({
+ startedAfter: NEXT_SNAPSHOT_AT,
+ countOnly: true,
+ }),
+ })
+ )
+ })
+
+ it('retains the rows and pending count when refreshing fails', async () => {
+ newCount = 2
+ const hook = renderHookWithClient(() => useLogs())
+ unmount = hook.unmount
+ await flushQueries()
+ const originalRows = hook.result().list.data
+ failRefresh = true
+ await act(async () => {
+ await expect(hook.result().list.refetch({ throwOnError: true })).rejects.toThrow(
+ 'Refresh failed'
+ )
+ })
+ await flushQueries()
+
+ expect(hook.result().list.data).toBe(originalRows)
+ expect(hook.result().count.data).toBe(2)
+ })
+
+ it('preserves resolved relative date bounds until explicit refresh', async () => {
+ vi.setSystemTime(new Date(SNAPSHOT_AT))
+ const hook = renderHookWithClient(() =>
+ useLogs({ ...LOG_FILTERS, timeRange: 'Past 30 minutes' })
+ )
+ unmount = hook.unmount
+ await flushQueries()
+ const firstStartDate = hook.result().list.data?.pages[0]?.query.startDate
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(10 * 60 * 1000)
+ await hook.result().list.fetchNextPage()
+ })
+ await flushQueries()
+
+ expect(hook.result().list.data?.pages[1]?.query.startDate).toBe(firstStartDate)
+ expect(mockRequestJson).toHaveBeenCalledWith(
+ expect.anything(),
+ expect.objectContaining({
+ query: expect.objectContaining({ cursor: 'next-page', startDate: firstStartDate }),
+ })
+ )
+
+ await act(async () => {
+ await hook.result().list.refetch()
+ })
+ await flushQueries()
+ expect(hook.result().list.data?.pages[0]?.query.startDate).not.toBe(firstStartDate)
+ expect(hook.result().list.data?.pages[1]?.query.startDate).toBe(
+ hook.result().list.data?.pages[0]?.query.startDate
+ )
+ })
+
+ it('stops pagination when a row moves across the sort cursor and recovers on refresh', async () => {
+ const hook = renderHookWithClient(() => useLogs({ ...LOG_FILTERS, sortBy: 'cost' }))
+ unmount = hook.unmount
+ await flushQueries()
+ const originalRows = hook.result().list.data?.pages[0]?.logs
+ revision = 'cost-changed-revision'
+ await act(async () => {
+ await hook.result().list.fetchNextPage()
+ })
+ await flushQueries()
+
+ expect(hook.result().list.data?.pages[0]?.logs).toBe(originalRows)
+ expect(hook.result().list.data?.pages[1]).toMatchObject({ logs: [], snapshotChanged: true })
+ expect(hook.result().list.hasNextPage).toBe(false)
+
+ snapshotAt = NEXT_SNAPSHOT_AT
+ await act(async () => {
+ await hook.result().list.refetch()
+ })
+ await flushQueries()
+ expect(hook.result().list.data?.pages.every((page) => !page.snapshotChanged)).toBe(true)
+ expect(hook.result().list.data?.pages[1]?.logs).toHaveLength(1)
+ })
+
+ it.each(['late-visible error', 'changed cost sort'])(
+ 'signals a %s without replacing rows or repeatedly scanning history',
+ async (change) => {
+ const filters: LogFilters = {
+ ...LOG_FILTERS,
+ ...(change === 'late-visible error' ? { level: 'error' } : { sortBy: 'cost' }),
+ }
+ const hook = renderHookWithClient(() => useLogs(filters))
+ unmount = hook.unmount
+ await flushQueries()
+ const originalRows = hook.result().list.data
+ expect(hook.result().updates.data).toBe(false)
+ revision = 'changed-revision'
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(LOG_SNAPSHOT_UPDATES_STALE_TIME)
+ })
+ await flushQueries()
+
+ expect(hook.result().count.data).toBe(0)
+ expect(hook.result().updates.data).toBe(true)
+ expect(hook.result().list.data).toBe(originalRows)
+ const revisionCalls = () =>
+ mockRequestJson.mock.calls.filter(
+ ([, { query }]) => query.includeRevision && query.countOnly
+ )
+ expect(revisionCalls()).toHaveLength(1)
+ expect(revisionCalls()[0][1].query).toMatchObject({
+ snapshotAt: SNAPSHOT_AT,
+ ...(filters.level === 'all' ? {} : { level: filters.level }),
+ sortBy: filters.sortBy,
+ })
+ expect(revisionCalls()[0][1].query.startedAfter).toBeUndefined()
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(LOG_SNAPSHOT_UPDATES_STALE_TIME * 2)
+ })
+ expect(revisionCalls()).toHaveLength(1)
+
+ snapshotAt = NEXT_SNAPSHOT_AT
+ await act(async () => {
+ await hook.result().list.refetch()
+ })
+ await flushQueries()
+ expect(hook.result().updates.data).toBe(false)
+ }
+ )
+
+ it('resets the indicator when filters change and checks the same filters as the list', async () => {
+ let filters = LOG_FILTERS
+ newCount = 8
+ const hook = renderHookWithClient(() => useLogs(filters))
+ unmount = hook.unmount
+ await flushQueries()
+ expect(hook.result().count.data).toBe(8)
+
+ filters = {
+ ...LOG_FILTERS,
+ workflowIds: ['wf-1'],
+ level: 'error',
+ sortBy: 'cost',
+ sortOrder: 'asc',
+ }
+ newCount = 0
+ hook.rerender()
+ expect(hook.result().count.data).toBeUndefined()
+ await flushQueries()
+
+ expect(hook.result().count.data).toBe(0)
+ expect(mockRequestJson).toHaveBeenCalledWith(
+ expect.anything(),
+ expect.objectContaining({
+ query: expect.objectContaining({
+ workspaceId: 'ws-1',
+ workflowIds: 'wf-1',
+ level: 'error',
+ sortBy: 'cost',
+ sortOrder: 'asc',
+ startedAfter: SNAPSHOT_AT,
+ countOnly: true,
+ }),
+ })
+ )
+ })
+
+ it('detects new logs after an empty initial list', async () => {
+ mockRequestJson.mockImplementation(async (_contract, { query }) => ({
+ data: [],
+ nextCursor: null,
+ ...(query.startedAfter ? { total: newCount } : { snapshotAt }),
+ }))
+ const hook = renderHookWithClient(() => useLogs())
+ unmount = hook.unmount
+ await flushQueries()
+ newCount = 1
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(NEW_LOG_COUNT_STALE_TIME)
+ })
+ await flushQueries()
+ expect(hook.result().list.data?.pages[0]?.logs).toEqual([])
+ expect(hook.result().count.data).toBe(1)
+ })
+
+ it('does not poll when the view is disabled', async () => {
+ const hook = renderHookWithClient(() => useLogs(LOG_FILTERS, false))
+ unmount = hook.unmount
+ await act(async () => {
+ await vi.advanceTimersByTimeAsync(NEW_LOG_COUNT_STALE_TIME * 2)
+ })
+ expect(mockRequestJson).not.toHaveBeenCalled()
+ })
+})
+
describe('useCancelExecution', () => {
beforeEach(() => {
vi.clearAllMocks()
@@ -66,6 +420,37 @@ describe('useCancelExecution', () => {
vi.useRealTimers()
})
+ it('reconciles the stopped row without replacing the displayed snapshot', async () => {
+ mockRequestJson.mockResolvedValueOnce({ success: true }).mockResolvedValueOnce({
+ data: { id: 'log-1', executionId: 'execution-1', status: 'cancelled' },
+ })
+ const hook = renderHookWithClient(() => useCancelExecution('workspace-1'))
+ const key = logKeys.snapshot('workspace-1', LOG_FILTERS)
+ hook.queryClient.setQueryData(key, {
+ pages: [
+ {
+ logs: [{ id: 'log-1', executionId: 'execution-1', status: 'running' }, { id: 'log-2' }],
+ snapshotAt: SNAPSHOT_AT,
+ nextCursor: 'older-page',
+ },
+ ],
+ pageParams: [null],
+ })
+ await act(async () => {
+ await hook.result().mutateAsync({ workflowId: 'wf-1', executionId: 'execution-1' })
+ })
+ expect(hook.queryClient.getQueryData(key)).toMatchObject({
+ pages: [
+ {
+ logs: [{ id: 'log-1', status: 'cancelled' }, { id: 'log-2' }],
+ snapshotAt: SNAPSHOT_AT,
+ nextCursor: 'older-page',
+ },
+ ],
+ })
+ hook.unmount()
+ })
+
it('polls the execution and keeps reconciling until it is terminal', async () => {
mockRequestJson
.mockResolvedValueOnce({
@@ -124,7 +509,24 @@ describe('useCancelExecution', () => {
})
.mockResolvedValueOnce({ data: { id: 'log-1', status } })
- const { result, unmount } = renderHookWithClient(() => useCancelExecution('workspace-1'))
+ const { result, unmount, queryClient } = renderHookWithClient(() =>
+ useCancelExecution('workspace-1')
+ )
+ const key = logKeys.snapshot('workspace-1', LOG_FILTERS)
+ queryClient.setQueryData(key, {
+ pages: [
+ {
+ logs: [{ id: 'log-1', executionId: 'execution-1', status: 'running' }],
+ snapshotAt: SNAPSHOT_AT,
+ nextCursor: null,
+ },
+ ],
+ pageParams: [null],
+ })
+ queryClient.setQueryData(logKeys.detail('workspace-1', 'log-1'), {
+ id: 'log-1',
+ status: 'running',
+ })
await act(async () => {
await expect(
@@ -134,6 +536,12 @@ describe('useCancelExecution', () => {
})
).rejects.toThrow(`Run finished as ${status} before cancellation was confirmed`)
})
+ expect(queryClient.getQueryData(key)).toMatchObject({
+ pages: [{ logs: [{ id: 'log-1', status }], snapshotAt: SNAPSHOT_AT }],
+ })
+ expect(queryClient.getQueryData(logKeys.detail('workspace-1', 'log-1'))).toMatchObject({
+ status,
+ })
unmount()
}
)
diff --git a/apps/sim/hooks/queries/logs.ts b/apps/sim/hooks/queries/logs.ts
index b1d357e2ed3..b741d2cc2e0 100644
--- a/apps/sim/hooks/queries/logs.ts
+++ b/apps/sim/hooks/queries/logs.ts
@@ -18,6 +18,7 @@ import {
getExecutionSnapshotContract,
getLogByExecutionIdContract,
getLogDetailContract,
+ type ListLogsQuery,
listLogsContract,
type WorkflowLogDetail,
type WorkflowLogSummary,
@@ -34,6 +35,9 @@ export type LogSortBy = 'date' | 'duration' | 'cost' | 'status'
export type LogSortOrder = 'asc' | 'desc'
export const LOG_LIST_STALE_TIME = 30 * 1000
+export const LOG_SNAPSHOT_STALE_TIME = 'static' as const
+export const NEW_LOG_COUNT_STALE_TIME = 10 * 1000
+export const LOG_SNAPSHOT_UPDATES_STALE_TIME = 60 * 1000
export const LOG_DETAIL_STALE_TIME = 30 * 1000
export const LOG_BY_EXECUTION_STALE_TIME = 30 * 1000
export const LOG_DASHBOARD_STATS_STALE_TIME = 30 * 1000
@@ -46,6 +50,20 @@ export const logKeys = {
lists: () => [...logKeys.all, 'list'] as const,
list: (workspaceId: string | undefined, filters: LogFilters) =>
[...logKeys.lists(), workspaceId ?? '', filters] as const,
+ snapshot: (workspaceId: string | undefined, filters: LogFilters) =>
+ [...logKeys.list(workspaceId, filters), 'snapshot'] as const,
+ newCounts: () => [...logKeys.all, 'newCount'] as const,
+ newCount: (workspaceId: string | undefined, filters: LogFilters, snapshotAt?: string) =>
+ [...logKeys.newCounts(), workspaceId ?? '', filters, snapshotAt ?? ''] as const,
+ updates: () => [...logKeys.all, 'updates'] as const,
+ update: (workspaceId: string | undefined, snapshot?: LogSnapshotPage) =>
+ [
+ ...logKeys.updates(),
+ workspaceId ?? '',
+ snapshot?.query,
+ snapshot?.snapshotAt,
+ snapshot?.revision,
+ ] as const,
details: () => [...logKeys.all, 'detail'] as const,
detail: (workspaceId: string | undefined, logId: string | undefined) =>
[...logKeys.details(), workspaceId ?? '', logId ?? ''] as const,
@@ -114,7 +132,11 @@ function applyFilterParams(
}
}
-function buildListQuery(workspaceId: string, filters: LogFilters, cursor: string | null) {
+function buildListQuery(
+ workspaceId: string,
+ filters: LogFilters,
+ cursor: string | null
+): ListLogsQuery {
const params = new URLSearchParams()
applyFilterParams(params, filters)
@@ -186,6 +208,126 @@ export function useLogsList(
})
}
+interface LogSnapshotCursor {
+ cursor: string
+ snapshotAt?: string
+ revision?: string
+ query: ListLogsQuery
+}
+
+interface LogSnapshotPage extends LogsPage {
+ snapshotAt?: string
+ revision?: string
+ query: ListLogsQuery
+ snapshotChanged: boolean
+}
+
+/** Keeps the displayed list stable until an explicit refresh, including across focus and invalidation. */
+export function useLogsSnapshot(
+ workspaceId: string | undefined,
+ filters: LogFilters,
+ options?: Pick
+) {
+ return useInfiniteQuery({
+ queryKey: logKeys.snapshot(workspaceId, filters),
+ queryFn: async ({ pageParam, signal }): Promise => {
+ const query = pageParam?.query ?? buildListQuery(workspaceId as string, filters, null)
+ const result = await requestJson(listLogsContract, {
+ query: {
+ ...query,
+ cursor: pageParam?.cursor,
+ snapshotAt: pageParam?.snapshotAt ?? 'now',
+ includeRevision: true,
+ },
+ signal,
+ })
+ const snapshotChanged = Boolean(pageParam && result.revision !== pageParam.revision)
+ return {
+ logs: snapshotChanged ? [] : result.data,
+ nextCursor: snapshotChanged ? null : result.nextCursor,
+ snapshotAt: result.snapshotAt,
+ revision: result.revision,
+ query,
+ snapshotChanged,
+ }
+ },
+ enabled: Boolean(workspaceId) && (options?.enabled ?? true),
+ staleTime: LOG_SNAPSHOT_STALE_TIME,
+ placeholderData: keepPreviousData,
+ initialPageParam: null as LogSnapshotCursor | null,
+ getNextPageParam: (lastPage): LogSnapshotCursor | undefined =>
+ lastPage.nextCursor
+ ? {
+ cursor: lastPage.nextCursor,
+ snapshotAt: lastPage.snapshotAt,
+ revision: lastPage.revision,
+ query: lastPage.query,
+ }
+ : undefined,
+ })
+}
+
+/** Checks older matching rows for late visibility or sort changes, stopping once refresh is needed. */
+export function useLogSnapshotUpdates(
+ workspaceId: string | undefined,
+ snapshot: LogSnapshotPage | undefined,
+ options?: Pick
+) {
+ return useQuery({
+ queryKey: logKeys.update(workspaceId, snapshot),
+ queryFn: async ({ signal }) => {
+ const result = await requestJson(listLogsContract, {
+ query: {
+ ...snapshot?.query,
+ workspaceId: workspaceId as string,
+ snapshotAt: snapshot?.snapshotAt,
+ countOnly: true,
+ includeRevision: true,
+ },
+ signal,
+ })
+ return result.revision !== snapshot?.revision
+ },
+ enabled: (query) =>
+ Boolean(workspaceId) &&
+ Boolean(snapshot?.revision) &&
+ (options?.enabled ?? true) &&
+ query.state.data !== true,
+ initialData: false,
+ staleTime: LOG_SNAPSHOT_UPDATES_STALE_TIME,
+ refetchInterval: (query) => (query.state.data ? false : LOG_SNAPSHOT_UPDATES_STALE_TIME),
+ refetchOnMount: false,
+ refetchOnWindowFocus: false,
+ refetchOnReconnect: false,
+ })
+}
+
+/** Polls only the count of new matching runs; it never writes to the displayed list. */
+export function useNewLogCount(
+ workspaceId: string | undefined,
+ filters: LogFilters,
+ snapshotAt: string | undefined,
+ options?: Pick
+) {
+ return useQuery({
+ queryKey: logKeys.newCount(workspaceId, filters, snapshotAt),
+ queryFn: async ({ signal }) => {
+ const result = await requestJson(listLogsContract, {
+ query: {
+ ...buildListQuery(workspaceId as string, filters, null),
+ countOnly: true,
+ startedAfter: snapshotAt,
+ },
+ signal,
+ })
+ return result.total ?? 0
+ },
+ enabled: Boolean(workspaceId) && Boolean(snapshotAt) && (options?.enabled ?? true),
+ staleTime: NEW_LOG_COUNT_STALE_TIME,
+ refetchInterval: NEW_LOG_COUNT_STALE_TIME,
+ })
+}
+
interface UseLogDetailOptions {
enabled?: boolean
refetchInterval?:
@@ -260,6 +402,30 @@ async function pollForTerminalExecution(params: {
if (detail) {
queryClient.setQueryData(logKeys.byExecution(workspaceId, executionId), detail)
queryClient.setQueryData(logKeys.detail(workspaceId, detail.id), detail)
+ const confirmedDetail = detail
+ /** Reconcile the acted-on row without inserting new runs into a manual snapshot. */
+ queryClient.setQueriesData>({ queryKey: logKeys.lists() }, (old) => {
+ if (!old) return old
+ return {
+ ...old,
+ pages: old.pages.map((page) => ({
+ ...page,
+ logs: page.logs.map((log) =>
+ log.executionId === executionId
+ ? {
+ ...log,
+ status: confirmedDetail.status,
+ level: confirmedDetail.level,
+ duration: confirmedDetail.duration,
+ cost: confirmedDetail.cost,
+ pauseSummary: confirmedDetail.pauseSummary,
+ hasPendingPause: confirmedDetail.hasPendingPause,
+ }
+ : log
+ ),
+ })),
+ }
+ })
}
await queryClient.invalidateQueries({ queryKey: logKeys.lists(), refetchType: 'active' })
@@ -410,14 +576,29 @@ export function useCancelExecution(workspaceId: string) {
return { previousQueries, affectedLogId, previousDetail }
},
- onError: (_err, _variables, context) => {
+ onError: (_err, { executionId }, context) => {
for (const [queryKey, data] of context?.previousQueries ?? []) {
- queryClient.setQueryData(queryKey, data)
+ const previousLog = data?.pages
+ .find((page) => page.logs.some((log) => log.executionId === executionId))
+ ?.logs.find((log) => log.executionId === executionId)
+ if (!previousLog) continue
+ queryClient.setQueryData>(queryKey, (current) => {
+ if (!current) return current
+ return {
+ ...current,
+ pages: current.pages.map((page) => ({
+ ...page,
+ logs: page.logs.map((log) =>
+ log.executionId === executionId && log.status === 'cancelling' ? previousLog : log
+ ),
+ })),
+ }
+ })
}
if (context?.affectedLogId && context.previousDetail !== undefined) {
- queryClient.setQueryData(
+ queryClient.setQueryData(
logKeys.detail(workspaceId, context.affectedLogId),
- context.previousDetail
+ (current) => (current?.status === 'cancelling' ? context.previousDetail : current)
)
}
},
diff --git a/apps/sim/lib/api/contracts/logs.ts b/apps/sim/lib/api/contracts/logs.ts
index 6e73aac9c03..f923b13a3db 100644
--- a/apps/sim/lib/api/contracts/logs.ts
+++ b/apps/sim/lib/api/contracts/logs.ts
@@ -1,5 +1,5 @@
import { z } from 'zod'
-import { userFileSchema } from '@/lib/api/contracts/primitives'
+import { booleanQueryFlagSchema, userFileSchema } from '@/lib/api/contracts/primitives'
import { defineRouteContract } from '@/lib/api/contracts/types'
const comparisonOperatorSchema = z.enum(['=', '>', '<', '>=', '<=', '!='])
@@ -40,8 +40,18 @@ export const listLogsQuerySchema = logFilterQuerySchema.extend({
sortOrder: logSortOrderSchema,
/** Also run a COUNT(*) under the same filters and return it as `total`. */
includeTotal: z.coerce.boolean().optional(),
+ /** Skip fetching and sorting log rows; return total and any requested revision. */
+ countOnly: booleanQueryFlagSchema.optional(),
+ /** Include a fingerprint of matching rows and their sort values for change detection. */
+ includeRevision: booleanQueryFlagSchema.optional(),
+ /** Bound run start times for pagination; mutable fields remain live on the server. */
+ snapshotAt: z.union([z.literal('now'), z.iso.datetime()]).optional(),
+ /** Count or list runs started after the displayed snapshot, within the other filters. */
+ startedAfter: z.iso.datetime().optional(),
})
+export type ListLogsQuery = z.input
+
export const logDetailQuerySchema = z.object({
workspaceId: z.string().min(1),
})
@@ -332,7 +342,10 @@ export type WorkflowLogRow = WorkflowLogSummary &
export const listLogsResponseSchema = z.object({
data: z.array(workflowLogSummarySchema),
nextCursor: z.string().nullable(),
- /** Total rows matching the filters; present only when `includeTotal` was set. */
+ /** Server-resolved upper bound for a manually refreshed list. */
+ snapshotAt: z.iso.datetime().optional(),
+ revision: z.string().max(160).optional(),
+ /** Total rows matching the filters; present when `includeTotal` or `countOnly` was set. */
total: z.number().optional(),
})
diff --git a/apps/sim/lib/logs/list-logs.test.ts b/apps/sim/lib/logs/list-logs.test.ts
index 3ab31ecc96e..b30f13cea91 100644
--- a/apps/sim/lib/logs/list-logs.test.ts
+++ b/apps/sim/lib/logs/list-logs.test.ts
@@ -44,9 +44,9 @@ vi.mock('@/lib/logs/folder-expansion', () => ({
expandFolderIdsWithDescendants: vi.fn(async (_ws: string, ids: string | undefined) => ids),
}))
-import type { ListLogsParams } from './list-logs'
-import { readLogs } from './list-logs'
-import { decodeLogSortCursor } from './sort-cursor'
+import { listLogsQuerySchema } from '@/lib/api/contracts/logs'
+import { type ReadLogsParams, readLogs } from '@/lib/logs/list-logs'
+import { decodeLogSortCursor } from '@/lib/logs/sort-cursor'
afterAll(resetDbChainMock)
@@ -99,14 +99,15 @@ function jobRow(overrides: Record = {}) {
}
}
-function baseParams(overrides: Partial = {}): ListLogsParams {
+function baseParams(overrides: Partial = {}): ReadLogsParams {
return {
workspaceId: 'ws-1',
limit: 100,
sortBy: 'date',
sortOrder: 'desc',
+ hideCostInfo: false,
...overrides,
- } as ListLogsParams
+ }
}
describe('readLogs', () => {
@@ -182,6 +183,128 @@ describe('readLogs', () => {
expect(result.data).toHaveLength(1)
expect(result.data[0].workflowId).toBe('wf-1')
})
+
+ it('resolves the snapshot on the server and applies its upper bound to both run sources', async () => {
+ const now = '2026-09-24T15:45:00.000Z'
+ vi.useFakeTimers()
+ vi.setSystemTime(new Date(now))
+ try {
+ const result = await readLogs(baseParams({ snapshotAt: 'now' }))
+ expect(result.snapshotAt).toBe(now)
+ for (const table of [workflowExecutionLogs, jobExecutionLogs]) {
+ expect(dbChainMockFns.where).toHaveBeenCalledWith(
+ expect.objectContaining({
+ args: expect.arrayContaining([{ type: 'lte', args: [table.startedAt, new Date(now)] }]),
+ })
+ )
+ }
+ } finally {
+ vi.useRealTimers()
+ }
+ })
+
+ it('counts new workflow and job runs with the same exclusive lower bound as the row query', async () => {
+ const startedAfter = '2026-09-24T15:45:00.000Z'
+ queueTableRows(workflowExecutionLogs, [])
+ queueTableRows(jobExecutionLogs, [])
+ queueTableRows(workflowExecutionLogs, [{ count: 3 }])
+ queueTableRows(jobExecutionLogs, [{ count: 2 }])
+
+ const result = await readLogs(baseParams({ startedAfter, includeTotal: true, limit: 1 }))
+ expect(result.total).toBe(5)
+ for (const table of [workflowExecutionLogs, jobExecutionLogs]) {
+ const matchingCalls = dbChainMockFns.where.mock.calls.filter(([condition]) =>
+ condition.args.some(
+ (item: { type: string; args: unknown[] }) =>
+ item.type === 'gt' &&
+ item.args[0] === table.startedAt &&
+ item.args[1] instanceof Date &&
+ item.args[1].toISOString() === startedAfter
+ )
+ )
+ expect(matchingCalls).toHaveLength(2)
+ }
+ })
+
+ it('validates snapshot boundaries without changing ordinary list requests', () => {
+ expect(listLogsQuerySchema.parse({ workspaceId: 'ws-1' }).snapshotAt).toBeUndefined()
+ for (const field of ['snapshotAt', 'startedAfter']) {
+ expect(
+ listLogsQuerySchema.safeParse({ workspaceId: 'ws-1', [field]: 'invalid' }).success
+ ).toBe(false)
+ }
+ })
+
+ it('counts matching new runs without fetching or sorting log rows', async () => {
+ const startedAfter = '2026-09-24T15:45:00.000Z'
+ queueTableRows(workflowExecutionLogs, [{ count: 3 }])
+ queueTableRows(jobExecutionLogs, [{ count: 2 }])
+
+ const result = await readLogs(baseParams({ countOnly: true, startedAfter, sortBy: 'cost' }))
+
+ expect(result).toEqual({ data: [], nextCursor: null, total: 5 })
+ expect(dbChainMockFns.select).toHaveBeenCalledTimes(2)
+ expect(dbChainMockFns.orderBy).not.toHaveBeenCalled()
+ expect(dbChainMockFns.limit).not.toHaveBeenCalled()
+ for (const table of [workflowExecutionLogs, jobExecutionLogs]) {
+ expect(dbChainMockFns.where).toHaveBeenCalledWith(
+ expect.objectContaining({
+ args: expect.arrayContaining([
+ { type: 'gt', args: [table.startedAt, new Date(startedAfter)] },
+ ]),
+ })
+ )
+ }
+ })
+
+ it('preserves workflow-specific filters for count-only requests', async () => {
+ queueTableRows(workflowExecutionLogs, [{ count: 3 }])
+
+ const result = await readLogs(baseParams({ countOnly: true, workflowIds: 'wf-1' }))
+
+ expect(result.total).toBe(3)
+ expect(dbChainMockFns.select).toHaveBeenCalledTimes(1)
+ expect(dbChainMockFns.orderBy).not.toHaveBeenCalled()
+ })
+
+ it('captures rows and the membership revision in one repeatable-read snapshot', async () => {
+ queueTableRows(workflowExecutionLogs, [workflowRow()])
+ queueTableRows(jobExecutionLogs, [])
+ queueTableRows(workflowExecutionLogs, [{ count: 1, revision: '1234567890123456789' }])
+ queueTableRows(jobExecutionLogs, [{ count: 0, revision: '0' }])
+
+ const result = await readLogs(baseParams({ includeRevision: true, snapshotAt: 'now' }))
+
+ expect(dbChainMockFns.transaction).toHaveBeenCalledWith(expect.any(Function), {
+ isolationLevel: 'repeatable read',
+ accessMode: 'read only',
+ })
+ expect(result.data).toHaveLength(1)
+ expect(result.revision).toBe('1:1234567890123456789:0:0')
+ expect(result.total).toBeUndefined()
+ })
+
+ it('reads a membership revision without loading rows or joining unrelated tables', async () => {
+ queueTableRows(workflowExecutionLogs, [{ count: 2, revision: '9' }])
+ queueTableRows(jobExecutionLogs, [{ count: 1, revision: '4' }])
+
+ const result = await readLogs(
+ baseParams({ countOnly: true, includeRevision: true, snapshotAt: '2026-01-01T00:00:00.000Z' })
+ )
+
+ expect(result.revision).toBe('2:9:1:4')
+ expect(result.total).toBe(3)
+ expect(result.data).toEqual([])
+ expect(dbChainMockFns.orderBy).not.toHaveBeenCalled()
+ expect(dbChainMockFns.leftJoin).not.toHaveBeenCalled()
+ })
+
+ it('parses the count-only query flag without treating false as true', () => {
+ expect(listLogsQuerySchema.parse({ workspaceId: 'ws-1', countOnly: true }).countOnly).toBe(true)
+ expect(listLogsQuerySchema.parse({ workspaceId: 'ws-1', countOnly: 'false' }).countOnly).toBe(
+ false
+ )
+ })
})
describe('readLogs cost projection', () => {
diff --git a/apps/sim/lib/logs/list-logs.ts b/apps/sim/lib/logs/list-logs.ts
index 2fd1baa5acb..c3919ad7e9b 100644
--- a/apps/sim/lib/logs/list-logs.ts
+++ b/apps/sim/lib/logs/list-logs.ts
@@ -67,18 +67,32 @@ type SortOrder = 'asc' | 'desc'
*/
export async function readLogs(params: ReadLogsParams): Promise {
params.signal?.throwIfAborted()
+ const folderIds = params.folderIds
+ ? await expandFolderIdsWithDescendants(params.workspaceId, params.folderIds)
+ : params.folderIds
+ const resolvedParams = { ...params, folderIds }
+ params.signal?.throwIfAborted()
+ if (params.includeRevision) {
+ return dbReplica.transaction((tx) => readLogsWithDatabase(resolvedParams, tx), {
+ isolationLevel: 'repeatable read',
+ accessMode: 'read only',
+ })
+ }
+ return readLogsWithDatabase(resolvedParams, dbReplica)
+}
+
+async function readLogsWithDatabase(
+ params: ReadLogsParams,
+ database: Pick
+): Promise {
+ params.signal?.throwIfAborted()
+ const snapshotAt = params.snapshotAt === 'now' ? new Date().toISOString() : params.snapshotAt
const { hideCostInfo } = params
const sortBy = params.sortBy as SortBy
const sortOrder = params.sortOrder as SortOrder
const cursor = params.cursor ? decodeLogSortCursor(params.cursor) : null
- // Expand selected folders to include descendants (matches the route behavior),
- // without mutating the caller's params object.
- const folderIds = params.folderIds
- ? await expandFolderIdsWithDescendants(params.workspaceId, params.folderIds)
- : params.folderIds
- params.signal?.throwIfAborted()
- const p: ReadLogsParams = { ...params, folderIds }
+ const p = params
const workflowSortExpr: SQL = (() => {
switch (sortBy) {
@@ -118,6 +132,12 @@ export async function readLogs(params: ReadLogsParams): Promise 0 && !levelList.some((l) => l === 'error' || l === 'info')
const includeJobLogs = !hasWorkflowSpecificFilters && !triggersExcludeJobs && !levelExcludesJobs
- const workflowQuery = dbReplica
- .select({
- id: workflowExecutionLogs.id,
- workflowId: workflowExecutionLogs.workflowId,
- executionId: workflowExecutionLogs.executionId,
- deploymentVersionId: workflowExecutionLogs.deploymentVersionId,
- level: workflowExecutionLogs.level,
- status: workflowExecutionLogs.status,
- trigger: workflowExecutionLogs.trigger,
- startedAt: workflowExecutionLogs.startedAt,
- endedAt: workflowExecutionLogs.endedAt,
- totalDurationMs: workflowExecutionLogs.totalDurationMs,
- costTotal: workflowExecutionLogs.costTotal,
- createdAt: workflowExecutionLogs.createdAt,
- workflowName: workflow.name,
- workflowDescription: workflow.description,
- workflowFolderId: workflow.folderId,
- workflowWorkspaceId: workflow.workspaceId,
- workflowCreatedAt: workflow.createdAt,
- workflowUpdatedAt: workflow.updatedAt,
- pausedStatus: pausedExecutions.status,
- pausedTotalPauseCount: pausedExecutions.totalPauseCount,
- pausedResumedCount: pausedExecutions.resumedCount,
- deploymentVersion: workflowDeploymentVersion.version,
- deploymentVersionName: workflowDeploymentVersion.name,
- executionOrigin: workflowExecutionOriginSql().as('execution_origin'),
- sortValue: sql`${workflowSortExpr}`.as('sort_value'),
- })
- .from(workflowExecutionLogs)
- .leftJoin(pausedExecutions, eq(pausedExecutions.executionId, workflowExecutionLogs.executionId))
- .leftJoin(
- workflowDeploymentVersion,
- eq(workflowDeploymentVersion.id, workflowExecutionLogs.deploymentVersionId)
- )
- .leftJoin(workflow, eq(workflowExecutionLogs.workflowId, workflow.id))
- .where(and(...workflowConditions))
- .orderBy(orderByClause(workflowSortExpr), dir(workflowExecutionLogs.id))
- .limit(fetchSize)
+ const workflowQuery = p.countOnly
+ ? Promise.resolve([])
+ : database
+ .select({
+ id: workflowExecutionLogs.id,
+ workflowId: workflowExecutionLogs.workflowId,
+ executionId: workflowExecutionLogs.executionId,
+ deploymentVersionId: workflowExecutionLogs.deploymentVersionId,
+ level: workflowExecutionLogs.level,
+ status: workflowExecutionLogs.status,
+ trigger: workflowExecutionLogs.trigger,
+ startedAt: workflowExecutionLogs.startedAt,
+ endedAt: workflowExecutionLogs.endedAt,
+ totalDurationMs: workflowExecutionLogs.totalDurationMs,
+ costTotal: workflowExecutionLogs.costTotal,
+ createdAt: workflowExecutionLogs.createdAt,
+ workflowName: workflow.name,
+ workflowDescription: workflow.description,
+ workflowFolderId: workflow.folderId,
+ workflowWorkspaceId: workflow.workspaceId,
+ workflowCreatedAt: workflow.createdAt,
+ workflowUpdatedAt: workflow.updatedAt,
+ pausedStatus: pausedExecutions.status,
+ pausedTotalPauseCount: pausedExecutions.totalPauseCount,
+ pausedResumedCount: pausedExecutions.resumedCount,
+ deploymentVersion: workflowDeploymentVersion.version,
+ deploymentVersionName: workflowDeploymentVersion.name,
+ executionOrigin: workflowExecutionOriginSql().as('execution_origin'),
+ sortValue: sql`${workflowSortExpr}`.as('sort_value'),
+ })
+ .from(workflowExecutionLogs)
+ .leftJoin(
+ pausedExecutions,
+ eq(pausedExecutions.executionId, workflowExecutionLogs.executionId)
+ )
+ .leftJoin(
+ workflowDeploymentVersion,
+ eq(workflowDeploymentVersion.id, workflowExecutionLogs.deploymentVersionId)
+ )
+ .leftJoin(workflow, eq(workflowExecutionLogs.workflowId, workflow.id))
+ .where(and(...workflowConditions))
+ .orderBy(orderByClause(workflowSortExpr), dir(workflowExecutionLogs.id))
+ .limit(fetchSize)
const jobConditions: SQL[] = [eq(jobExecutionLogs.workspaceId, p.workspaceId)]
let jobFilterConditions: SQL[] = jobConditions
if (includeJobLogs) {
+ if (snapshotAt) {
+ jobConditions.push(lte(jobExecutionLogs.startedAt, new Date(snapshotAt)))
+ }
+ if (p.startedAfter) {
+ jobConditions.push(gt(jobExecutionLogs.startedAt, new Date(p.startedAfter)))
+ }
if (p.level && p.level !== 'all') {
const levels = p.level.split(',').filter(Boolean)
const jobLevelConditions: SQL[] = []
@@ -302,27 +333,28 @@ export async function readLogs(params: ReadLogsParams): Promise`${jobExecutionLogs.executionData}->'trigger'->>'source'`,
- sortValue: sql`${jobSortExpr}`.as('sort_value'),
- })
- .from(jobExecutionLogs)
- .where(and(...jobConditions))
- .orderBy(orderByClause(jobSortExpr), dir(jobExecutionLogs.id))
- .limit(fetchSize)
- : Promise.resolve([])
+ const jobQuery =
+ includeJobLogs && !p.countOnly
+ ? database
+ .select({
+ id: jobExecutionLogs.id,
+ executionId: jobExecutionLogs.executionId,
+ level: jobExecutionLogs.level,
+ status: jobExecutionLogs.status,
+ trigger: jobExecutionLogs.trigger,
+ startedAt: jobExecutionLogs.startedAt,
+ endedAt: jobExecutionLogs.endedAt,
+ totalDurationMs: jobExecutionLogs.totalDurationMs,
+ cost: jobExecutionLogs.cost,
+ createdAt: jobExecutionLogs.createdAt,
+ jobTitle: sql`${jobExecutionLogs.executionData}->'trigger'->>'source'`,
+ sortValue: sql`${jobSortExpr}`.as('sort_value'),
+ })
+ .from(jobExecutionLogs)
+ .where(and(...jobConditions))
+ .orderBy(orderByClause(jobSortExpr), dir(jobExecutionLogs.id))
+ .limit(fetchSize)
+ : Promise.resolve([])
const [workflowRows, jobRows] = await Promise.all([workflowQuery, jobQuery])
params.signal?.throwIfAborted()
@@ -447,35 +479,64 @@ export async function readLogs(params: ReadLogsParams): Promise`COUNT(*)` })
+ let revision: string | undefined
+ if (p.includeTotal || p.countOnly || p.includeRevision) {
+ const workflowRevisionValue =
+ sortBy === 'date'
+ ? sql`${workflowExecutionLogs.id}`
+ : sql`concat_ws(':', ${workflowExecutionLogs.id}, ${workflowSortExpr})`
+ const jobRevisionValue =
+ sortBy === 'date'
+ ? sql`${jobExecutionLogs.id}`
+ : sql`concat_ws(':', ${jobExecutionLogs.id}, ${jobSortExpr})`
+ const workflowCountBase = database
+ .select({
+ count: sql`COUNT(*)`,
+ revision: p.includeRevision
+ ? sql`COALESCE(bit_xor(hashtextextended(${workflowRevisionValue}, 0)), 0)::text`
+ : sql`NULL`,
+ })
.from(workflowExecutionLogs)
- .leftJoin(
- pausedExecutions,
- eq(pausedExecutions.executionId, workflowExecutionLogs.executionId)
- )
- .leftJoin(
- workflowDeploymentVersion,
- eq(workflowDeploymentVersion.id, workflowExecutionLogs.deploymentVersionId)
- )
- .leftJoin(workflow, eq(workflowExecutionLogs.workflowId, workflow.id))
- .where(and(...workflowFilterConditions))
+ const workflowCountWithPauses = levelList.includes('pending')
+ ? workflowCountBase.leftJoin(
+ pausedExecutions,
+ eq(pausedExecutions.executionId, workflowExecutionLogs.executionId)
+ )
+ : workflowCountBase
+ const workflowCountWithFilters = hasWorkflowSpecificFilters
+ ? workflowCountWithPauses.leftJoin(
+ workflow,
+ eq(workflowExecutionLogs.workflowId, workflow.id)
+ )
+ : workflowCountWithPauses
+ const workflowCountQuery = workflowCountWithFilters.where(and(...workflowFilterConditions))
const jobCountQuery = includeJobLogs
- ? dbReplica
- .select({ count: sql`COUNT(*)` })
+ ? database
+ .select({
+ count: sql`COUNT(*)`,
+ revision: p.includeRevision
+ ? sql`COALESCE(bit_xor(hashtextextended(${jobRevisionValue}, 0)), 0)::text`
+ : sql`NULL`,
+ })
.from(jobExecutionLogs)
.where(and(...jobFilterConditions))
- : Promise.resolve([{ count: 0 }])
+ : Promise.resolve([{ count: 0, revision: '0' }])
const [workflowCount, jobCount] = await Promise.all([workflowCountQuery, jobCountQuery])
params.signal?.throwIfAborted()
- total = Number(workflowCount[0]?.count ?? 0) + Number(jobCount[0]?.count ?? 0)
+ const workflowTotal = Number(workflowCount[0]?.count ?? 0)
+ const jobTotal = Number(jobCount[0]?.count ?? 0)
+ if (p.includeTotal || p.countOnly) total = workflowTotal + jobTotal
+ if (p.includeRevision) {
+ revision = `${workflowTotal}:${workflowCount[0]?.revision ?? '0'}:${jobTotal}:${jobCount[0]?.revision ?? '0'}`
+ }
}
params.signal?.throwIfAborted()
return {
data: page.map((row) => row.summary),
nextCursor,
+ ...(snapshotAt ? { snapshotAt } : {}),
+ ...(revision !== undefined ? { revision } : {}),
...(total !== undefined ? { total } : {}),
}
}