diff --git a/apps/desktop/e2e/background-executor.spec.ts b/apps/desktop/e2e/background-executor.spec.ts index 5eef9b4e582..06cc2dbcfd7 100644 --- a/apps/desktop/e2e/background-executor.spec.ts +++ b/apps/desktop/e2e/background-executor.spec.ts @@ -384,7 +384,10 @@ test.describe('background executor', () => { const deviceId = await registeredDevice(sim) const pulls = () => sim.requests.filter((request) => request.startsWith('GET /api/desktop/inbox')).length + const claimAttempts = () => + sim.requests.filter((request) => request === 'POST /api/desktop/tool/claim').length const pullsBefore = pulls() + const claimAttemptsBefore = claimAttempts() const gated = sim.issue( deviceId, @@ -399,7 +402,8 @@ test.describe('background executor', () => { await check('F: the declined call is never claimed, and its command never runs', async () => { await sleep(RECONCILE_MS * 2) - expect(sim.requireCall(gated).claims).toBe(0) + // Attempts, not successful claims: the fixture refuses a declined call's claim itself. + expect(claimAttempts() - claimAttemptsBefore).toBe(0) expect(sim.requireCall(gated).completions).toEqual([]) expect(readFileSafe(marker)).toBe('') }) diff --git a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts index be9dd5041d7..6a4fe359933 100644 --- a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts +++ b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts @@ -48,6 +48,7 @@ const ARRIVAL_MS = 60_000 /** First requests to a route compile it, which takes minutes on a cold dev app. */ const COMPILE_MS = 300_000 const REGISTRATION_PATH = '/api/desktop/devices' +const DESKTOP_COMPLETION_PATH = '/api/desktop/tool/complete' /** Sim's answer to registration on an install that cannot run the background executor. */ const executorUnavailable = (answer: Record) => { answer.enabled = false @@ -120,6 +121,7 @@ test.describe('desktop tools against a live Sim', () => { await app?.close().catch(() => {}) app = undefined proxy.rewriteAnswer(REGISTRATION_PATH, executorUnavailable) + proxy.rewriteAnswer(DESKTOP_COMPLETION_PATH, undefined) rmSync(scratch, { recursive: true, force: true }) }) @@ -377,7 +379,7 @@ test.describe('desktop tools against a live Sim', () => { entry.method === 'POST' && entry.path === '/api/mothership/chat' /** The background executor's report of a call's result. */ const isDesktopCompletion = (method: string, path: string) => - method === 'POST' && path === '/api/desktop/tool/complete' + method === 'POST' && path === DESKTOP_COMPLETION_PATH /** A client tool's report of its own result. */ const isToolReport = (method: string, path: string) => method === 'POST' && path === '/api/copilot/confirm' @@ -830,6 +832,7 @@ test.describe('desktop tools against a live Sim', () => { // Reaches Sim on release although the cut made the app give up on it, as a report already on // the wire would: the app cannot know it landed, so it reports again once back online. const completion = proxy.hold(isDesktopCompletion, { deliverIfAbandoned: true }) + const answers = proxy.recordAnswers(DESKTOP_COMPLETION_PATH) await send(page, '[network-cut] read my notes') await completion.arrival(ARRIVAL_MS, 'The result report') @@ -837,7 +840,8 @@ test.describe('desktop tools against a live Sim', () => { await expect.poll(() => completion.isAbandoned, { timeout: 15_000 }).toBe(true) completion.release() await expect.poll(() => callState(chatId), { timeout: 30_000 }).toMatch(/^completed/) - await sleep(5_000) + // The late report alone resumed the agent while the app was still offline. + await agent.waitForResume(() => Boolean(agent.resultFor(callId)), 60_000) const restoredAt = Date.now() proxy.restoreNetwork() @@ -852,8 +856,10 @@ test.describe('desktop tools against a live Sim', () => { await expect.poll(() => retried().length, { timeout: 60_000 }).toBeGreaterThan(0) for (const entry of retried()) expect(entry.status).toBeLessThan(300) - await agent.waitForResume(() => Boolean(agent.resultFor(callId)), 60_000) await proxy.settled(30_000) + // Sim had already recorded the late report, so every retry the app heard back from is a no-op. + expect(answers.length).toBeGreaterThan(0) + expect(answers).toEqual(answers.map(() => expect.objectContaining({ outcome: 'duplicate' }))) const delivered = agent.resumes.filter((resume) => resume.results.some((entry) => entry.callId === callId) ) diff --git a/apps/desktop/e2e/fixtures/live-sim.ts b/apps/desktop/e2e/fixtures/live-sim.ts index a27873c9ac2..ed3ff12f055 100644 --- a/apps/desktop/e2e/fixtures/live-sim.ts +++ b/apps/desktop/e2e/fixtures/live-sim.ts @@ -259,6 +259,18 @@ export class SimProxy { else this.answerRewrites.delete(path) } + /** + * Sim's JSON answers to `path` that reach the app from now on, recorded as they pass and left + * unchanged. An answer to a request whose client gave up reaches no one, so it is not recorded. + */ + recordAnswers(path: string): Record[] { + const answers: Record[] = [] + this.rewriteAnswer(path, (body) => { + answers.push(structuredClone(body)) + }) + return answers + } + /** * Cuts the network between the app and Sim: every open connection drops mid-flight, the * doorbell stream included, and new ones are reset until `restoreNetwork`. diff --git a/apps/desktop/src/main/terminal/index.ts b/apps/desktop/src/main/terminal/index.ts index 20ee810f9ee..c9334dab154 100644 --- a/apps/desktop/src/main/terminal/index.ts +++ b/apps/desktop/src/main/terminal/index.ts @@ -1420,7 +1420,7 @@ export class TerminalService { const screen = (await session.readScrollback(STARTUP_SCREEN_LINES)).output.trim() throw new TerminalError( 'NO_SHELL_INTEGRATION', - `The shell has been running its startup files for over ${SHELL_STARTUP_BOUNDS.startingMs / 1000} s without reaching a prompt, so nothing was run. Its screen:\n${screen || '(empty)'}\nIf it is waiting for an answer, ask the user to answer it in that terminal (terminalId ${session.terminalId}), then run the command again.` + `The shell began its startup files over ${SHELL_STARTUP_BOUNDS.startingMs / 1000} s ago but never reached a prompt Sim can track, so nothing was run. Its screen:\n${screen || '(empty)'}\nIf a startup file is waiting for an answer, ask the user to answer it in that terminal (terminalId ${session.terminalId}), then run the command again. If the screen shows a prompt, a startup file replaced the shell (such as exec tmux or exec fish), so ask the user to run the command themselves.` ) } case 'not-instrumented': diff --git a/apps/desktop/src/main/terminal/session.ts b/apps/desktop/src/main/terminal/session.ts index 6c64e85d3af..60afc3e54e3 100644 --- a/apps/desktop/src/main/terminal/session.ts +++ b/apps/desktop/src/main/terminal/session.ts @@ -328,7 +328,6 @@ export class TerminalSession { private columns: number private lines: number private shellIntegration = false - /** The shell has begun the startup files we generated, so its integration is on the way. */ /** When the shell began the startup files we generated, or null while it has not. */ private startupBegunAt: number | null = null private readonly spawnedAt = Date.now() diff --git a/apps/desktop/src/main/terminal/shell-integration.test.ts b/apps/desktop/src/main/terminal/shell-integration.test.ts index 573f601484a..9c16588b5f0 100644 --- a/apps/desktop/src/main/terminal/shell-integration.test.ts +++ b/apps/desktop/src/main/terminal/shell-integration.test.ts @@ -1,5 +1,5 @@ import { spawnSync } from 'node:child_process' -import { existsSync, mkdtempSync, writeFileSync } from 'node:fs' +import { existsSync, mkdirSync, mkdtempSync, writeFileSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' import { describe, expect, it } from 'vitest' @@ -140,6 +140,37 @@ describe('startup marker', () => { } ) + it.skipIf(!existsSync('/bin/zsh'))( + "zsh keeps integrating when the user's .zshenv moves ZDOTDIR, and runs all their files there", + () => { + const home = mkdtempSync(join(tmpdir(), 'sim-shell-home-')) + const userDir = join(home, '.config', 'zsh') + mkdirSync(userDir, { recursive: true }) + writeFileSync(join(home, '.zshenv'), 'export ZDOTDIR="$HOME/.config/zsh"\n') + writeFileSync(join(userDir, '.zprofile'), 'USER_PROFILE_RAN=1\n') + writeFileSync(join(userDir, 'plugins.zsh'), 'USER_PLUGIN_RAN=1\n') + writeFileSync(join(userDir, '.zshrc'), 'USER_RC_RAN=1\nsource "$ZDOTDIR/plugins.zsh"\n') + writeFileSync(join(userDir, '.zlogin'), 'echo "login=$USER_RC_RAN"\n') + const env = { PATH: process.env.PATH ?? '/usr/bin:/bin', HOME: home } + const launch = buildShellLaunch('zsh', mkdtempSync(join(tmpdir(), 'sim-zsh-')), NONCE, env) + + const output = spawnSync( + '/bin/zsh', + [ + ...launch.args, + '-i', + '-c', + 'echo "profile=$USER_PROFILE_RAN rc=$USER_RC_RAN plugin=$USER_PLUGIN_RAN zdotdir=$ZDOTDIR"; whence -w __sim_precmd', + ], + { env: { ...env, ...launch.env }, encoding: 'utf8' } + ).stdout + + expect(output).toContain(`profile=1 rc=1 plugin=1 zdotdir=${userDir}`) + expect(output).toContain('__sim_precmd: function') + expect(output).toContain('login=1') + } + ) + it.skipIf(!existsSync('/bin/bash'))("bash sends it before the user's files", () => { const home = userHome('.bashrc') const env = { PATH: process.env.PATH ?? '/usr/bin:/bin', HOME: home } diff --git a/apps/desktop/src/main/terminal/shell-integration.ts b/apps/desktop/src/main/terminal/shell-integration.ts index b0f770f430e..da7349bbb27 100644 --- a/apps/desktop/src/main/terminal/shell-integration.ts +++ b/apps/desktop/src/main/terminal/shell-integration.ts @@ -166,8 +166,17 @@ function findTerminator(buffer: string, from: number): { index: number; length: * prompt, aliases, and PATH win over ours. */ function writeZshFiles(dir: string, nonce: string, originalZdotdir: string): void { - const sourceOriginal = (file: string) => - `[ -f "$SIM_ZDOTDIR_ORIG/${file}" ] && builtin source "$SIM_ZDOTDIR_ORIG/${file}"` + // The user's file runs with ZDOTDIR at their own directory, so it finds its siblings and + // plugins there. A file that moves ZDOTDIR (an XDG `.zshenv`) moves where the rest of theirs are + // read from; ZDOTDIR then points back here, so zsh reads our next file. + const sourceOriginal = (file: string) => `if [ -f "$SIM_ZDOTDIR_ORIG/${file}" ]; then + __sim_zdotdir="$ZDOTDIR" + ZDOTDIR="$SIM_ZDOTDIR_ORIG" + builtin source "$SIM_ZDOTDIR_ORIG/${file}" + SIM_ZDOTDIR_ORIG="\${ZDOTDIR:-$HOME}" + ZDOTDIR="$__sim_zdotdir" + builtin unset __sim_zdotdir +fi` // `.zshenv` is the first file zsh reads, so the startup marker goes out before any of the user's // files run. Only from an interactive shell: a script's output must not carry it. diff --git a/apps/desktop/src/main/terminal/shell-startup.test.ts b/apps/desktop/src/main/terminal/shell-startup.test.ts index 4df9be581ab..f734aad5ec9 100644 --- a/apps/desktop/src/main/terminal/shell-startup.test.ts +++ b/apps/desktop/src/main/terminal/shell-startup.test.ts @@ -9,6 +9,7 @@ import { TerminalService } from '@/main/terminal' */ const pty = vi.hoisted(() => ({ emit: null as ((data: string) => void) | null, + exit: null as (() => void) | null, writes: [] as string[], })) @@ -18,7 +19,9 @@ vi.mock('@lydell/node-pty', () => ({ onData: (handler: (data: string) => void) => { pty.emit = handler }, - onExit: () => {}, + onExit: (handler: () => void) => { + pty.exit = handler + }, write: (data: string) => pty.writes.push(data), resize: vi.fn(), kill: vi.fn(), @@ -71,6 +74,7 @@ beforeEach(() => { afterEach(() => { vi.useRealTimers() pty.emit = null + pty.exit = null pty.writes.length = 0 vi.unstubAllEnvs() }) @@ -122,6 +126,20 @@ describe('a shell that is still starting', () => { expect(response).toMatchObject({ ok: false, code: 'NO_SHELL_INTEGRATION' }) expect(response.error).toContain('[oh-my-zsh] Would you like to update? [Y/n]') expect(response.error).toContain('ask the user to answer it') + expect(response.error).toContain('a startup file replaced the shell') + expect(pty.writes).toEqual([]) + terminal.dispose() + }) + + it('reports the session closed when the shell exits while its startup files run', async () => { + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-exit', 'run', { command: 'echo hi' }) + shell(STARTUP) + await vi.advanceTimersByTimeAsync(2_000) + pty.exit?.() + const response = await running + + expect(response).toMatchObject({ ok: false, code: 'SESSION_CLOSED' }) expect(pty.writes).toEqual([]) terminal.dispose() }) diff --git a/apps/sim/app/api/mothership/events/route.test.ts b/apps/sim/app/api/mothership/events/route.test.ts index c410c711850..bb83ece059f 100644 --- a/apps/sim/app/api/mothership/events/route.test.ts +++ b/apps/sim/app/api/mothership/events/route.test.ts @@ -220,4 +220,25 @@ describe('Mothership owner-scoped event stream', () => { expect(authorize).not.toHaveBeenCalled() expect(authMockFns.mockGetSession).toHaveBeenCalledTimes(1) }) + + it('tells each member whether a workspace chat is their own, never whose it is', async () => { + const abort = new AbortController() + const response = await GET(request('workspaceId=ws-1', abort.signal)) + if (!response.body) throw new Error('The event stream has no body') + const chunks: string[] = [] + const collected = collect(response.body, chunks) + await vi.advanceTimersByTimeAsync(0) + emit({ workspaceId: 'ws-1', userId: 'user-1', chatId: 'own-chat', type: 'started' }) + emit({ workspaceId: 'ws-1', userId: 'teammate-1', chatId: 'teammate-chat', type: 'started' }) + emit({ workspaceId: 'ws-1', chatId: 'unknown-owner-chat', type: 'started' }) + abort.abort() + await collected + const payloads = chunks.map((chunk) => JSON.parse(chunk.split('data: ')[1])) + expect(payloads).toEqual([ + expect.objectContaining({ chatId: 'own-chat', ownChat: true }), + expect.objectContaining({ chatId: 'teammate-chat', ownChat: false }), + expect.not.objectContaining({ ownChat: expect.anything() }), + ]) + expect(chunks.join('')).not.toMatch(/userId|teammate-1|user-1/) + }) }) diff --git a/apps/sim/app/api/mothership/events/route.ts b/apps/sim/app/api/mothership/events/route.ts index 4b0e5695c3e..deb78a641b3 100644 --- a/apps/sim/app/api/mothership/events/route.ts +++ b/apps/sim/app/api/mothership/events/route.ts @@ -30,7 +30,7 @@ const mothershipEventsHandler = createWorkspaceSSE({ label: 'mothership-events', subscriptions: [ { - subscribe: (workspaceId, send) => { + subscribe: (workspaceId, send, viewerUserId) => { if (!chatPubSub) return () => {} return chatPubSub.onStatusChanged((event) => { if (event.workspaceId !== workspaceId) return @@ -38,6 +38,8 @@ const mothershipEventsHandler = createWorkspaceSSE({ chatId: event.chatId, type: event.type, ...(event.streamId ? { streamId: event.streamId } : {}), + // Whether the chat is the viewer's own, never whose it is; absent when unknown. + ...(event.userId ? { ownChat: event.userId === viewerUserId } : {}), timestamp: Date.now(), }) }) diff --git a/apps/sim/app/workspace/[workspaceId]/layout.test.tsx b/apps/sim/app/workspace/[workspaceId]/layout.test.tsx index 6cdd0319416..f42e3b26788 100644 --- a/apps/sim/app/workspace/[workspaceId]/layout.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/layout.test.tsx @@ -10,6 +10,7 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' const { mockBrandingProvider, mockIsDesktopPresenceAvailable, + mockHasSignedInDesktopExecutor, mockWorkspaceChrome, mockGetOrgWhitelabelSettings, mockPrefetchWorkspaceHostContext, @@ -19,6 +20,7 @@ const { } = vi.hoisted(() => ({ mockBrandingProvider: vi.fn(({ children }: { children: ReactNode }) => children), mockIsDesktopPresenceAvailable: vi.fn(() => false), + mockHasSignedInDesktopExecutor: vi.fn(async () => false), mockWorkspaceChrome: vi.fn( ({ children }: { children: ReactNode; sidebar: ReactNode }) => children ), @@ -77,6 +79,10 @@ vi.mock('@/lib/desktop/executor/presence', () => ({ isDesktopPresenceAvailable: mockIsDesktopPresenceAvailable, })) +vi.mock('@/lib/desktop/executor/repository', () => ({ + hasSignedInDesktopExecutor: mockHasSignedInDesktopExecutor, +})) + vi.mock('@/app/workspace/[workspaceId]/w/components/sidebar/sidebar', () => ({ Sidebar: () => null, })) @@ -213,20 +219,35 @@ describe('WorkspaceLayout host context', () => { expect(mockGetOrgWhitelabelSettings).not.toHaveBeenCalled() }) + async function sidebarProps() { + mockWorkspaceChrome.mockClear() + const element = await WorkspaceLayout({ + children:
Workspace child
, + params: Promise.resolve({ workspaceId: 'workspace-b' }), + }) + renderToStaticMarkup(element) + return mockWorkspaceChrome.mock.calls[0][0].sidebar + } + + it('tells the sidebar the desktop executor runs only where presence is tracked', async () => { + mockIsDesktopPresenceAvailable.mockReturnValue(false) + mockHasSignedInDesktopExecutor.mockResolvedValue(true) + + expect(await sidebarProps()).toMatchObject({ + props: { desktopExecutor: { available: false, registered: false } }, + }) + mockHasSignedInDesktopExecutor.mockResolvedValue(false) + }) + it.each([true, false])( - 'tells the sidebar the desktop executor runs only where presence is tracked (%s)', - async (presenceAvailable) => { - mockIsDesktopPresenceAvailable.mockReturnValue(presenceAvailable) - mockWorkspaceChrome.mockClear() - - const element = await WorkspaceLayout({ - children:
Workspace child
, - params: Promise.resolve({ workspaceId: 'workspace-b' }), - }) - renderToStaticMarkup(element) + 'tells the sidebar whether the viewer has a desktop that runs their turns (%s)', + async (registered) => { + mockIsDesktopPresenceAvailable.mockReturnValue(true) + mockHasSignedInDesktopExecutor.mockResolvedValueOnce(registered) - const { sidebar } = mockWorkspaceChrome.mock.calls[0][0] - expect(sidebar).toMatchObject({ props: { desktopExecutorAvailable: presenceAvailable } }) + expect(await sidebarProps()).toMatchObject({ + props: { desktopExecutor: { available: true, registered } }, + }) } ) }) diff --git a/apps/sim/app/workspace/[workspaceId]/layout.tsx b/apps/sim/app/workspace/[workspaceId]/layout.tsx index fe983f71dc1..55c5f662861 100644 --- a/apps/sim/app/workspace/[workspaceId]/layout.tsx +++ b/apps/sim/app/workspace/[workspaceId]/layout.tsx @@ -5,7 +5,10 @@ import { SettingsNavigationProvider } from '@/components/settings/settings-navig import { getSession } from '@/lib/auth' import { getActiveOrganizationId } from '@/lib/auth/session-response' import { isDashboardsEnabled } from '@/lib/dashboards/feature-flag' -import { isDesktopBackgroundExecutorAvailable } from '@/lib/desktop/executor/availability' +import { + hasDesktopBackgroundExecutor, + isDesktopBackgroundExecutorAvailable, +} from '@/lib/desktop/executor/availability' import { isMothershipModelSelectorEnabled, isPlanModeEnabled } from '@/lib/mothership/feature-flags' import { resolveOrganizationEntryPath } from '@/lib/navigation/resolve-app-entry' import { isTableRowTtlEnabled } from '@/lib/table/ttl-availability' @@ -68,6 +71,7 @@ export default async function WorkspaceLayout({ planModeEnabled, organizationHref, dashboardsEnabled, + desktopExecutorRegistered, ] = await Promise.all([ cookies(), hostContext.hostOrganizationId @@ -85,6 +89,7 @@ export default async function WorkspaceLayout({ isPlanModeEnabled(), resolveOrganizationEntryPath(session), isDashboardsEnabled(hostContext.hostOrganizationId), + hasDesktopBackgroundExecutor(session.user.id), prefetchWorkspaceAccess(queryClient, workspaceId, principal), prefetchWorkspaceForkAvailability(queryClient, workspaceId, principal, hostContext), ]) @@ -122,7 +127,10 @@ export default async function WorkspaceLayout({ sidebar={ } initialSidebarCollapsed={initialSidebarCollapsed} diff --git a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx index 8010fb27410..44f17083375 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx @@ -124,7 +124,7 @@ import { useImportWorkflow } from '@/app/workspace/[workspaceId]/w/hooks' import { useCustomBlockOverlayVersion } from '@/blocks/custom/client-overlay' import { useWorkspaceAccessRequestFeatures } from '@/ee/access-requests/components/permission-access-boundary' import { useWorkspaceCredentials } from '@/hooks/queries/credentials' -import { useDesktopActivity } from '@/hooks/queries/desktop-activity' +import { useDesktopActivity, watchesDesktopActivity } from '@/hooks/queries/desktop-activity' import { useFolderMap, useFolders } from '@/hooks/queries/folders' import { type LogFilters, useLogsList } from '@/hooks/queries/logs' import type { MothershipChatMetadata } from '@/hooks/queries/mothership-chats' @@ -372,8 +372,11 @@ const DRAG_EXEMPT_CLASS = '[-webkit-app-region:no-drag]' interface SidebarProps { organizationHref: string | null - /** Whether this install runs the desktop background executor, so chats show desktop activity. */ - desktopExecutorAvailable: boolean + /** + * Whether this install runs the desktop background executor, and whether the user has a desktop + * registered to run their turns, so chats show desktop activity. + */ + desktopExecutor: { available: boolean; registered: boolean } } /** @@ -392,10 +395,7 @@ interface SidebarProps { * * @returns Sidebar with workflows panel */ -export const Sidebar = memo(function Sidebar({ - organizationHref, - desktopExecutorAvailable, -}: SidebarProps) { +export const Sidebar = memo(function Sidebar({ organizationHref, desktopExecutor }: SidebarProps) { const { isCollapsed: isCollapsedProp, isPeeking } = useSidebarChrome() const isCollapsed = isCollapsedProp && !isPeeking const params = useParams() @@ -890,14 +890,15 @@ export const Sidebar = memo(function Sidebar({ { enabled: chatEnabled && !permissionConfig.hideCopilot } ) + const desktopActivityWatched = watchesDesktopActivity(desktopExecutor) useMothershipChatEvents( workspaceId, chatEnabled && !permissionConfig.hideCopilot, - desktopExecutorAvailable + desktopActivityWatched ) - const { data: desktopActivity } = useDesktopActivity( + const desktopActivity = useDesktopActivity( workspaceId, - desktopExecutorAvailable && chatEnabled && !permissionConfig.hideCopilot + desktopActivityWatched && chatEnabled && !permissionConfig.hideCopilot ) const desktopActivityByChat = useMemo( () => new Map((desktopActivity ?? []).map((activity) => [activity.chatId, activity])), diff --git a/apps/sim/hooks/queries/desktop-activity.test.tsx b/apps/sim/hooks/queries/desktop-activity.test.tsx new file mode 100644 index 00000000000..dfb5c78d83b --- /dev/null +++ b/apps/sim/hooks/queries/desktop-activity.test.tsx @@ -0,0 +1,120 @@ +/** @vitest-environment jsdom */ + +import { act } from 'react' +import { jsonResponse } from '@sim/testing/helpers/http' +import { setupGlobalFetchMock } from '@sim/testing/mocks/fetch.mock' +import { QueryClient, QueryClientProvider } from '@tanstack/react-query' +import { createRoot, type Root } from 'react-dom/client' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +/** The page's real requests land on this fetch, installed fresh for each test. */ +let fetchMock = setupGlobalFetchMock() + +/** Requests the page sent for desktop activity. */ +function activityRequests(): number { + return fetchMock.mock.calls.filter(([input]) => String(input).includes('/api/desktop/activity')) + .length +} + +import { useDesktopActivity, watchesDesktopActivity } from '@/hooks/queries/desktop-activity' + +interface ActivityWatcherProps { + registered: boolean +} + +/** What the page last saw from the hook. */ +let shown: unknown + +function ActivityWatcher({ registered }: ActivityWatcherProps) { + shown = useDesktopActivity('ws-1', watchesDesktopActivity({ available: true, registered })) + return null +} + +let root: Root | undefined + +async function mount(registered: boolean) { + const container = document.createElement('div') + document.body.appendChild(container) + root = createRoot(container) + await act(async () => { + root?.render( + + + + ) + }) +} + +/** Long enough for several idle polls. */ +async function waitTwoMinutes() { + await act(async () => { + await vi.advanceTimersByTimeAsync(2 * 60 * 1000) + }) +} + +describe('desktop activity polling', () => { + beforeEach(() => { + vi.useFakeTimers() + vi.stubGlobal('IS_REACT_ACT_ENVIRONMENT', true) + fetchMock = setupGlobalFetchMock() + fetchMock.mockImplementation(async () => jsonResponse({ chats: [] })) + }) + + afterEach(() => { + act(() => root?.unmount()) + root = undefined + vi.useRealTimers() + vi.unstubAllGlobals() + Reflect.deleteProperty(window, 'simDesktop') + }) + + it('shows nothing once the page stops watching, though the query kept its last result', async () => { + const activity = [ + { chatId: 'chat-1', streamId: 'stream-1', state: 'running', deviceName: 'Work laptop' }, + ] + fetchMock.mockImplementationOnce(async () => jsonResponse({ chats: activity })) + const container = document.createElement('div') + document.body.appendChild(container) + root = createRoot(container) + const client = new QueryClient() + const render = (registered: boolean) => + act(async () => { + root?.render( + + + + ) + }) + + await render(true) + await act(async () => { + await vi.advanceTimersByTimeAsync(0) + }) + expect(shown).toEqual(activity) + + await render(false) + expect(shown).toBeUndefined() + }) + + it('never asks for a user without a desktop in a browser tab', async () => { + await mount(false) + await waitTwoMinutes() + + expect(activityRequests()).toBe(0) + }) + + it('polls for a user with a registered desktop', async () => { + await mount(true) + await waitTwoMinutes() + + expect(activityRequests()).toBeGreaterThanOrEqual(4) + }) + + it('polls in the desktop app before its first registration reaches the page', async () => { + Object.defineProperty(window, 'simDesktop', { value: {}, configurable: true }) + await mount(false) + await waitTwoMinutes() + + expect(activityRequests()).toBeGreaterThanOrEqual(4) + }) +}) diff --git a/apps/sim/hooks/queries/desktop-activity.ts b/apps/sim/hooks/queries/desktop-activity.ts index 494f2de62ae..4143a5602fe 100644 --- a/apps/sim/hooks/queries/desktop-activity.ts +++ b/apps/sim/hooks/queries/desktop-activity.ts @@ -6,13 +6,14 @@ import { type DesktopChatActivity, listDesktopActivityContract, } from '@/lib/api/contracts/desktop-executor' +import { isDesktopApp } from '@/lib/desktop' import { desktopActivityKeys } from '@/hooks/queries/utils/desktop-activity-keys' const DESKTOP_ACTIVITY_STALE_TIME = 10 * 1000 /** * Presence, approvals and new background turns change without a chat event this query hears, so - * it is re-read on a timer: often while a desktop runs a chat, rarely otherwise. Only installs that - * run the executor ever read it. + * it is re-read on a timer: often while a desktop runs a chat, rarely otherwise. Only users with a + * desktop ever read it. */ const DESKTOP_ACTIVITY_ACTIVE_REFETCH_MS = 15 * 1000 const DESKTOP_ACTIVITY_IDLE_REFETCH_MS = 30 * 1000 @@ -25,9 +26,29 @@ async function fetchDesktopActivity( return data.chats } -/** The user's chats in this workspace whose turn runs on one of their desktops. */ -export function useDesktopActivity(workspaceId: string | undefined, enabled: boolean) { - return useQuery({ +/** + * Whether a page shows background desktop activity: only for a user whose turns can run on one of + * their desktops. The page learns that when it loads, so the desktop app's own window also asks, + * since its first registration can land after the page rendered; a browser tab open across that + * first registration shows activity once reloaded. + */ +export function watchesDesktopActivity(desktopExecutor: { + available: boolean + registered: boolean +}): boolean { + return desktopExecutor.available && (desktopExecutor.registered || isDesktopApp()) +} + +/** + * The user's chats in this workspace whose turn runs on one of their desktops, or nothing while + * the page does not watch them: a disabled query keeps its last result, which must not outlive + * the reason it was read. + */ +export function useDesktopActivity( + workspaceId: string | undefined, + enabled: boolean +): DesktopChatActivity[] | undefined { + const { data } = useQuery({ queryKey: desktopActivityKeys.list(workspaceId), queryFn: ({ signal }) => fetchDesktopActivity(workspaceId as string, signal), enabled: Boolean(workspaceId) && enabled, @@ -38,4 +59,5 @@ export function useDesktopActivity(workspaceId: string | undefined, enabled: boo : DESKTOP_ACTIVITY_IDLE_REFETCH_MS, placeholderData: keepPreviousData, }) + return enabled ? data : undefined } diff --git a/apps/sim/hooks/use-mothership-chat-events.test.ts b/apps/sim/hooks/use-mothership-chat-events.test.ts index f8e58df7b90..a44c6d57a5a 100644 --- a/apps/sim/hooks/use-mothership-chat-events.test.ts +++ b/apps/sim/hooks/use-mothership-chat-events.test.ts @@ -391,6 +391,35 @@ describe('reflectBackgroundChatStatus', () => { expect(activityStale()).toBe(true) }) + it.each(['started', 'completed'])( + "leaves the desktop activity alone when a teammate's turn has %s", + (type) => { + showing('/workspace/ws-1/home', false) + + reflectBackgroundChatStatus( + queryClient, + 'ws-1', + JSON.stringify({ chatId: 'teammate-chat', type, streamId: 's-9', ownChat: false }), + true + ) + + expect(activityStale()).toBe(false) + } + ) + + it('refreshes the desktop activity when the viewer’s own turn starts', () => { + showing('/workspace/ws-1/home', false) + + reflectBackgroundChatStatus( + queryClient, + 'ws-1', + JSON.stringify({ chatId: 'chat-new', type: 'started', streamId: 's-4', ownChat: true }), + true + ) + + expect(activityStale()).toBe(true) + }) + it('leaves the desktop activity alone for a rename', () => { showing('/workspace/ws-1/home', false) diff --git a/apps/sim/hooks/use-mothership-chat-events.ts b/apps/sim/hooks/use-mothership-chat-events.ts index 02166feff8b..604082f5639 100644 --- a/apps/sim/hooks/use-mothership-chat-events.ts +++ b/apps/sim/hooks/use-mothership-chat-events.ts @@ -37,6 +37,8 @@ interface ChatStatusEventPayload { chatId?: string type?: ChatStatusEventType streamId?: string + /** Whether the chat is the viewer's own; absent when the server could not tell. */ + ownChat?: boolean } function isChatStatusEventType(value: unknown): value is ChatStatusEventType { @@ -102,6 +104,7 @@ function parseChatStatusEventPayload(data: unknown): ChatStatusEventPayload | nu ...(typeof record.chatId === 'string' ? { chatId: record.chatId } : {}), ...(isChatStatusEventType(record.type) ? { type: record.type } : {}), ...(typeof record.streamId === 'string' ? { streamId: record.streamId } : {}), + ...(typeof record.ownChat === 'boolean' ? { ownChat: record.ownChat } : {}), } } @@ -209,6 +212,8 @@ export function reflectBackgroundChatStatus( ): void { const payload = parseChatStatusEventPayload(data) if (payload?.type !== 'started' && payload?.type !== 'completed') return + // A teammate's turn never runs on this user's desktop, so it changes nothing here. + if (payload.ownChat === false) return // Read before the refresh below drops it. The events carry every member's chats in the // workspace; only this very turn, seen running on this user's own desktop, is theirs to be told // about. Matching the turn, not just the chat, keeps an earlier desktop turn from vouching for a diff --git a/apps/sim/lib/desktop/application/activity.integration.ts b/apps/sim/lib/desktop/application/activity.integration.ts index 855013e53da..37328eaa26b 100644 --- a/apps/sim/lib/desktop/application/activity.integration.ts +++ b/apps/sim/lib/desktop/application/activity.integration.ts @@ -24,10 +24,11 @@ import { workspace, } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { inArray } from 'drizzle-orm' +import { eq, inArray } from 'drizzle-orm' import { closeRedisConnection } from '@/lib/core/config/redis' import { listDesktopActivity } from '@/lib/desktop/application/activity' import { openDesktopInboxStream, registerDesktopDevice } from '@/lib/desktop/application/executor' +import { hasDesktopBackgroundExecutor } from '@/lib/desktop/executor/availability' import { isDesktopPresent } from '@/lib/desktop/executor/presence' import { createRunSegment } from '@/lib/mothership/async-runs/repository' @@ -256,4 +257,31 @@ describe.runIf(Boolean(redisUrl))('background desktop activity', () => { }) expect(otherWorkspace).toEqual([]) }) + + it('watches activity only for a user with a signed-in desktop that runs their turns', async () => { + const desktop = await signedInDesktop() + expect(await hasDesktopBackgroundExecutor(desktop.userId)).toBe(true) + + const withoutExecutor = await signedInDesktop() + await db + .update(desktopDevices) + .set({ capabilities: { browser: true } }) + .where(eq(desktopDevices.id, withoutExecutor.deviceId)) + const revoked = await signedInDesktop() + await db + .update(desktopDevices) + .set({ revokedAt: new Date() }) + .where(eq(desktopDevices.id, revoked.deviceId)) + const signedOut = await signedInDesktop() + await db.delete(session).where(eq(session.id, signedOut.principal.sessionId)) + const expired = await signedInDesktop() + await db + .update(session) + .set({ expiresAt: new Date(Date.now() - 60_000) }) + .where(eq(session.id, expired.principal.sessionId)) + + for (const owner of [withoutExecutor, revoked, signedOut, expired]) { + expect(await hasDesktopBackgroundExecutor(owner.userId)).toBe(false) + } + }) }) diff --git a/apps/sim/lib/desktop/application/activity.ts b/apps/sim/lib/desktop/application/activity.ts index 06e6b5dd674..6fe7439878a 100644 --- a/apps/sim/lib/desktop/application/activity.ts +++ b/apps/sim/lib/desktop/application/activity.ts @@ -38,8 +38,8 @@ async function readPresence(deviceId: string): Promise { /** * Which of the caller's chats in a workspace are running on one of their desktops, and in what * state. Lists only the caller's own runs, never their content, so it needs no workspace role. - * The endpoint is not flag-gated (runs already bound are listed until they end); the sidebar asks - * only while the background executor is on for the user. + * Runs already bound are listed until they end, whatever the install's executor availability; the + * sidebar asks only for a user with a desktop that runs their turns. */ export const listDesktopActivity = defineAuthorizedCredentialUserUseCase({ // permission-group-exempt: reports only the caller's own runs, with no content. diff --git a/apps/sim/lib/desktop/executor/availability.ts b/apps/sim/lib/desktop/executor/availability.ts index 59ef02e81f3..b49f852d8ed 100644 --- a/apps/sim/lib/desktop/executor/availability.ts +++ b/apps/sim/lib/desktop/executor/availability.ts @@ -1,4 +1,5 @@ import { isDesktopPresenceAvailable } from '@/lib/desktop/executor/presence' +import { hasSignedInDesktopExecutor } from '@/lib/desktop/executor/repository' /** * Whether new turns may bind to a desktop, and whether the workspace shows background desktop @@ -8,3 +9,12 @@ import { isDesktopPresenceAvailable } from '@/lib/desktop/executor/presence' export function isDesktopBackgroundExecutorAvailable(): boolean { return isDesktopPresenceAvailable() } + +/** + * Whether the user's turns can run on one of their desktops, so their workspace has background + * desktop activity to show. Most users never install the app, and their pages never ask. + */ +export async function hasDesktopBackgroundExecutor(userId: string): Promise { + if (!isDesktopBackgroundExecutorAvailable()) return false + return hasSignedInDesktopExecutor(userId) +} diff --git a/apps/sim/lib/desktop/executor/repository.ts b/apps/sim/lib/desktop/executor/repository.ts index 324e55f2de2..3b6df2918f1 100644 --- a/apps/sim/lib/desktop/executor/repository.ts +++ b/apps/sim/lib/desktop/executor/repository.ts @@ -5,12 +5,14 @@ import { copilotChats, copilotRuns, desktopDevices, + session, } from '@sim/db/schema' import { TERMINAL_TOOL_NAME } from '@sim/terminal-protocol' import { and, asc, eq, + gt, inArray, isNotNull, isNull, @@ -142,6 +144,11 @@ export async function upsertDesktopDevice(input: DesktopDeviceRegistration): Pro return Boolean(row) } +/** A device that registered a background executor a turn can be bound to. */ +function registersExecutor() { + return sql`coalesce((${desktopDevices.capabilities} ->> 'executor')::int, 0) >= 1` +} + export interface DesktopDeviceIdentity { deviceId: string userId: string @@ -165,15 +172,34 @@ export async function getBoundDesktopDevice( eq(desktopDevices.userId, identity.userId), eq(desktopDevices.sessionId, identity.sessionId), isNull(desktopDevices.revokedAt), - options.executor - ? sql`coalesce((${desktopDevices.capabilities} ->> 'executor')::int, 0) >= 1` - : undefined + options.executor ? registersExecutor() : undefined ) ) .limit(1) return row ?? null } +/** + * Whether the user has a desktop that registered a background executor and is still signed in: + * its session exists and has not expired. + */ +export async function hasSignedInDesktopExecutor(userId: string): Promise { + const [row] = await db + .select({ id: desktopDevices.id }) + .from(desktopDevices) + .innerJoin(session, eq(session.id, desktopDevices.sessionId)) + .where( + and( + eq(desktopDevices.userId, userId), + isNull(desktopDevices.revokedAt), + gt(session.expiresAt, new Date()), + registersExecutor() + ) + ) + .limit(1) + return Boolean(row) +} + /** Display and support only; written at most once a minute per device. */ export async function touchDesktopDevice(deviceId: string): Promise { await db diff --git a/apps/sim/lib/events/sse-endpoint.ts b/apps/sim/lib/events/sse-endpoint.ts index 0b5afb6a260..e9bb1a03955 100644 --- a/apps/sim/lib/events/sse-endpoint.ts +++ b/apps/sim/lib/events/sse-endpoint.ts @@ -18,7 +18,9 @@ import { getUserEntityPermissions } from '@/lib/workspaces/permissions/utils' interface SSESubscription { subscribe( workspaceId: string, - send: (eventName: string, data: Record) => void + send: (eventName: string, data: Record) => void, + /** The member the stream belongs to. */ + viewerUserId: string ): () => void /** Settles once the subscription receives events; the stream is announced only after it. */ ready?: () => Promise @@ -95,7 +97,7 @@ export function createWorkspaceSSE(config: WorkspaceSSEConfig) { return createSSEStream(request, { label: `${config.label}:workspace:${workspaceId}`, subscriptions: config.subscriptions.map((subscription) => ({ - subscribe: (send) => subscription.subscribe(workspaceId, send), + subscribe: (send) => subscription.subscribe(workspaceId, send, userId), ready: subscription.ready, })), }) diff --git a/apps/sim/lib/mothership/chat-status.test.ts b/apps/sim/lib/mothership/chat-status.test.ts index 567a10b33c5..a59c9a6f053 100644 --- a/apps/sim/lib/mothership/chat-status.test.ts +++ b/apps/sim/lib/mothership/chat-status.test.ts @@ -1,19 +1,35 @@ -import { describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { + type ChatStatusEvent, + chatPubSub, + publishChatStatusChanged, +} from '@/lib/mothership/chat-status' -const { publish } = vi.hoisted(() => ({ publish: vi.fn() })) -vi.mock('@/lib/events/pubsub', () => ({ - createPubSubChannel: () => ({ publish, subscribe: vi.fn(), dispose: vi.fn() }), -})) +/** Events as a subscriber receives them from the channel (process-local without Redis). */ +let received: ChatStatusEvent[] = [] +let unsubscribe: () => void = () => {} -import { publishChatStatusChanged } from '@/lib/mothership/chat-status' +beforeEach(() => { + received = [] + unsubscribe = chatPubSub?.onStatusChanged((event) => received.push(event)) ?? (() => {}) +}) + +afterEach(() => unsubscribe()) describe('chat status ownership', () => { - it('preserves the workspace event shape', () => { + it('names the chat owner on workspace events, so each member can tell their own chats', () => { publishChatStatusChanged( { workspaceId: 'ws-1', userId: 'user-1' }, { chatId: 'chat-1', type: 'renamed' } ) - expect(publish).toHaveBeenCalledWith({ workspaceId: 'ws-1', chatId: 'chat-1', type: 'renamed' }) + expect(received).toEqual([ + { workspaceId: 'ws-1', userId: 'user-1', chatId: 'chat-1', type: 'renamed' }, + ]) + }) + + it('publishes a workspace event without an owner when the publisher does not know it', () => { + publishChatStatusChanged({ workspaceId: 'ws-1' }, { chatId: 'chat-1', type: 'renamed' }) + expect(received).toEqual([{ workspaceId: 'ws-1', chatId: 'chat-1', type: 'renamed' }]) }) it.each(['completed'] as const)( @@ -23,12 +39,9 @@ describe('chat status ownership', () => { { organizationId: 'org-1', userId: 'user-1' }, { chatId: 'chat-1', type } ) - expect(publish).toHaveBeenCalledWith({ - organizationId: 'org-1', - userId: 'user-1', - chatId: 'chat-1', - type, - }) + expect(received).toEqual([ + { organizationId: 'org-1', userId: 'user-1', chatId: 'chat-1', type }, + ]) } ) @@ -36,7 +49,7 @@ describe('chat status ownership', () => { expect(() => publishChatStatusChanged({ organizationId: 'org-1' }, { chatId: 'chat-1', type: 'created' }) ).toThrow('Invalid organization chat owner') - expect(publish).not.toHaveBeenCalled() + expect(received).toEqual([]) }) it('refuses ambiguous workspace and organization ownership', () => { @@ -46,6 +59,6 @@ describe('chat status ownership', () => { { chatId: 'chat-1', type: 'created' } ) ).toThrow('Invalid organization chat owner') - expect(publish).not.toHaveBeenCalled() + expect(received).toEqual([]) }) }) diff --git a/apps/sim/lib/mothership/chat-status.ts b/apps/sim/lib/mothership/chat-status.ts index ac0b0a6fe03..94e830ca214 100644 --- a/apps/sim/lib/mothership/chat-status.ts +++ b/apps/sim/lib/mothership/chat-status.ts @@ -11,8 +11,12 @@ import { createPubSubChannel, type PubSubChannel } from '@/lib/events/pubsub' +/** + * A workspace event names the chat's owner when the publisher knows it, so each member's stream + * can tell their own chats from teammates'. Events from older pods carry no owner. + */ export type ChatStatusOwner = - | { workspaceId: string; organizationId?: never; userId?: never } + | { workspaceId: string; userId?: string; organizationId?: never } | { organizationId: string; userId: string; workspaceId?: never } export type ChatStatusEvent = ChatStatusOwner & { @@ -58,6 +62,10 @@ export function publishChatStatusChanged( ...event, }) } else if (chat.workspaceId) { - chatPubSub?.publishStatusChanged({ workspaceId: chat.workspaceId, ...event }) + chatPubSub?.publishStatusChanged({ + workspaceId: chat.workspaceId, + ...(chat.userId ? { userId: chat.userId } : {}), + ...event, + }) } } diff --git a/apps/sim/lib/mothership/tasks/wake.ts b/apps/sim/lib/mothership/tasks/wake.ts index a22750ef567..63f7614bdb2 100644 --- a/apps/sim/lib/mothership/tasks/wake.ts +++ b/apps/sim/lib/mothership/tasks/wake.ts @@ -36,7 +36,7 @@ const logger = createLogger('CopilotTaskWake') export async function runWakeTurn(input: WakeRequest): Promise { const { taskId, chatId, workspaceId, organizationId, userId, message, runId } = input const userMessageId = runId - const owner = organizationId ? { organizationId, userId } : { workspaceId: workspaceId! } + const owner = organizationId ? { organizationId, userId } : { workspaceId: workspaceId!, userId } chatPubSub?.publishStatusChanged({ ...owner, chatId,