diff --git a/src/main/__tests__/integration-event-bridge.test.ts b/src/main/__tests__/integration-event-bridge.test.ts index 7d6a5cc1..c5d8e517 100644 --- a/src/main/__tests__/integration-event-bridge.test.ts +++ b/src/main/__tests__/integration-event-bridge.test.ts @@ -77,6 +77,7 @@ function makeHarness(agents = ['alice', 'bob']): { const bridge = new IntegrationEventBridge({ getWorkspaceHandle: async () => ({ workspaceId: 'workspace-id', + localMountWorkspaceId: 'workspace-id', client: () => ({ subscribe(globs, onChange, options) { subscribeCalls.push({ globs: [...globs], onChange, options }) @@ -152,6 +153,24 @@ test('channel notification targets do not fall back to all project agents', asyn assert.deepEqual(harness.listAgentsCalls, []) }) +test('offline notification agents fall back to current project agents', async () => { + const harness = makeHarness(['alice', 'bob']) + + await harness.bridge.reconcile('project-1', [ + integration({ + provider: 'slack', + integrationId: 'slack-1', + mountPaths: ['/slack/channels'], + scope: { notifyAgents: ['claude-1'] } + }) + ]) + + await harness.emit(changeEvent('/slack/channels/general/messages/123.json', 'slack')) + + assert.deepEqual(harness.sent.map((message) => message.input.to), ['alice', 'bob']) + assert.deepEqual(harness.listAgentsCalls, ['project-1']) +}) + test('integration events watch selected relayfile mount paths', async () => { const harness = makeHarness() const slackIntegration = integration({ @@ -168,9 +187,11 @@ test('integration events watch selected relayfile mount paths', async () => { await harness.bridge.reconcile('project-1', [slackIntegration]) assert.deepEqual(harness.subscribeCalls[0].globs, [ + '/slack/channels/C123ABC/**', '/slack/channels/C123ABC__proj-cloud/**' ]) assert.deepEqual(integrationSubscriptionSummaries([slackIntegration])[0].watches, [ + '.integrations/slack/channels/C123ABC/**', '.integrations/slack/channels/C123ABC__proj-cloud/**' ]) @@ -179,6 +200,32 @@ test('integration events watch selected relayfile mount paths', async () => { assert.deepEqual(harness.sent.map((message) => message.input.to), ['alice']) assert.match(harness.sent[0].input.text, /Path: \.integrations\/slack\/channels\/C123ABC__proj-cloud\/messages\/1713220123_001100\/meta\.json/u) assert.match(harness.sent[0].input.text, /Relayfile path: \/slack\/channels\/C123ABC__proj-cloud\/messages\/1713220123_001100\/meta\.json/u) + + harness.sent.splice(0) + await harness.emit(changeEvent('/slack/channels/C123ABC/messages/1713220124_001100/meta.json', 'slack')) + assert.deepEqual(harness.sent.map((message) => message.input.to), ['alice']) +}) + +test('integration events preserve discovery mount paths', async () => { + const harness = makeHarness() + const slackIntegration = integration({ + provider: 'slack', + integrationId: 'slack-1', + mountPaths: ['/discovery/slack'] + }) + + await harness.bridge.reconcile('project-1', [slackIntegration]) + + assert.deepEqual(harness.subscribeCalls[0].globs, [ + '/discovery/slack/**' + ]) + assert.deepEqual(integrationSubscriptionSummaries([slackIntegration])[0].watches, [ + '.integrations/discovery/slack/**' + ]) + + await harness.emit(changeEvent('/discovery/slack/actions/create-message/.schema.json', 'slack')) + assert.deepEqual(harness.sent, []) + assert.deepEqual(harness.listAgentsCalls, []) }) test('resource alias mount paths inject the same relative event only once', async () => { diff --git a/src/main/__tests__/integration-remote-paths.test.ts b/src/main/__tests__/integration-remote-paths.test.ts new file mode 100644 index 00000000..d8509625 --- /dev/null +++ b/src/main/__tests__/integration-remote-paths.test.ts @@ -0,0 +1,64 @@ +import assert from 'node:assert/strict' +import { test } from 'node:test' + +import { + canShowRemoteDirectoryEntryForMountPaths, + canListRemoteDirectoryForMountPaths, + normalizeRemoteDirectoryPath, + remotePathName +} from '../integration-remote-paths.ts' + +test('remote directory paths reject traversal segments', () => { + assert.equal(normalizeRemoteDirectoryPath('/slack/channels'), '/slack/channels') + assert.equal(normalizeRemoteDirectoryPath('/slack/../channels'), null) + assert.equal(normalizeRemoteDirectoryPath('/slack/./channels'), null) + assert.equal(remotePathName('/slack/channels/C123'), 'C123') +}) + +test('remote directory listing is limited to configured mount roots', () => { + assert.equal(canListRemoteDirectoryForMountPaths('/slack/channels/C123', [ + '/slack/channels/C123' + ]), true) + assert.equal(canListRemoteDirectoryForMountPaths('/slack/channels/C123/messages', [ + '/slack/channels/C123' + ]), true) + assert.equal(canListRemoteDirectoryForMountPaths('/slack/channels', [ + '/slack/channels/C123' + ]), true) + assert.equal(canListRemoteDirectoryForMountPaths('/slack', [ + '/slack/channels/C123' + ]), false) + assert.equal(canListRemoteDirectoryForMountPaths('/slack/channels/C999', [ + '/slack/channels/C123' + ]), false) +}) + +test('remote directory listing permits provider discovery only for that provider', () => { + assert.equal(canListRemoteDirectoryForMountPaths('/discovery/slack/actions', [ + '/discovery/slack' + ]), true) + assert.equal(canListRemoteDirectoryForMountPaths('/discovery', [ + '/discovery/slack' + ]), false) + assert.equal(canListRemoteDirectoryForMountPaths('/discovery/github/actions', [ + '/discovery/slack' + ]), false) +}) + +test('remote directory entries are filtered to configured mount roots', () => { + assert.equal(canShowRemoteDirectoryEntryForMountPaths('/slack/channels/C123', [ + '/slack/channels/C123' + ]), true) + assert.equal(canShowRemoteDirectoryEntryForMountPaths('/slack/channels/C123/messages', [ + '/slack/channels/C123' + ]), true) + assert.equal(canShowRemoteDirectoryEntryForMountPaths('/slack/channels/C999', [ + '/slack/channels/C123' + ]), false) + assert.equal(canShowRemoteDirectoryEntryForMountPaths('/discovery', [ + '/discovery/slack' + ]), true) + assert.equal(canShowRemoteDirectoryEntryForMountPaths('/discovery/github', [ + '/discovery/slack' + ]), false) +}) diff --git a/src/main/integration-event-bridge.ts b/src/main/integration-event-bridge.ts index 661c8782..0f64a99d 100644 --- a/src/main/integration-event-bridge.ts +++ b/src/main/integration-event-bridge.ts @@ -1,6 +1,7 @@ import { watch, type FSWatcher } from 'node:fs' -import { stat } from 'node:fs/promises' -import { join, relative, resolve, sep } from 'node:path' +import { appendFile, mkdir, stat } from 'node:fs/promises' +import { homedir } from 'node:os' +import { dirname, join, relative, resolve, sep } from 'node:path' import { RelayFileClient, RelayFileSync, @@ -15,6 +16,7 @@ const INTEGRATION_EVENT_AGENT_NAME = 'pear-integration-events' const INTEGRATION_EVENT_SCOPES = ['relayfile:fs:read:/**'] const PROJECT_INTEGRATIONS_LINK_NAME = '.integrations' const RECENT_INJECTION_TTL_MS = 10_000 +const INTEGRATION_EVENT_LOG_PATH = join(homedir(), '.agentworkforce', 'pear', 'integration-events.log') type WatchRegistration = { glob: string @@ -68,6 +70,7 @@ type RelayfileEventClient = { type RelayfileWorkspaceHandle = { workspaceId: string + localMountWorkspaceId: string client(): RelayfileEventClient } @@ -81,7 +84,7 @@ type IntegrationEventBridgeDeps = { type EventWorkspaceHandleCache = { apiUrl: string accountKey: string - workspaceId: string + accountWorkspaceId: string handle: RelayfileWorkspaceHandle } @@ -145,13 +148,24 @@ function workspaceIdFromJwt(token: string | undefined): string | null { function canonicalMountPaths(integration: ConnectedIntegration): string[] { const provider = toRelayfileProvider(integration.provider) const mountPaths = integration.mountPaths.map((path) => { + const discovery = path.match(/^\/discovery(?:\/.*)?$/) + if (discovery) return path const prefixed = path.match(/^\/integrations\/[^/]+(\/.*)?$/) if (prefixed) return `/${provider}${prefixed[1] ?? ''}` const rootLevel = path.match(/^\/[^/]+(\/.*)?$/) if (rootLevel) return `/${provider}${rootLevel[1] ?? ''}` return path }) - return dedupeStrings(mountPaths) + return dedupeStrings([ + ...mountPaths, + ...mountPaths.map((path) => slackChannelIdFallbackMountPath(provider, path)).filter((path): path is string => path !== null) + ]) +} + +function slackChannelIdFallbackMountPath(provider: string, path: string): string | null { + if (provider !== 'slack') return null + const match = path.match(/^\/slack\/channels\/([^/_][^/]*)__[^/]+$/u) + return match?.[1] ? `/slack/channels/${match[1]}` : null } function watchGlobForPath(path: string): string { @@ -335,7 +349,12 @@ function createWorkspaceScopedEventClient( if (!active) return const changeEvent = filesystemEventToChangeEvent(client, workspaceId, event) Promise.resolve(onChange(changeEvent)).catch((error) => { - console.warn('[integration-events] Change handler failed:', toErrorMessage(error)) + logIntegrationEvent('change handler failed', { + workspaceId, + eventId: event.eventId, + path: event.path, + error: toErrorMessage(error) + }) }) } @@ -361,29 +380,42 @@ function createWorkspaceScopedEventClient( .then((token) => { if (!active) return const tokenWorkspaceId = workspaceIdFromJwt(token) - if (tokenWorkspaceId !== workspaceId) { - logIntegrationEvent('skipping remote stream without workspace JWT', { + if (tokenWorkspaceId && tokenWorkspaceId !== workspaceId) { + logIntegrationEvent('skipping remote stream with mismatched workspace JWT', { workspaceId, tokenWorkspaceId }) return } + logIntegrationEvent('remote stream starting', { + workspaceId, + globs + }) sync = new RelayFileSync({ client, workspaceId, token, onPollingFallback: (info) => { - console.warn('[integration-events] Relayfile event stream using polling fallback:', info.reason) + logIntegrationEvent('remote stream polling fallback', { + workspaceId, + reason: info.reason + }) } }) sync.on('event', handleEvent) sync.on('error', (error) => { - console.warn('[integration-events] Relayfile event stream error:', toErrorMessage(error)) + logIntegrationEvent('remote stream error', { + workspaceId, + error: toErrorMessage(error) + }) }) sync.start() }) .catch((error) => { - console.warn('[integration-events] Relayfile event stream token check failed:', toErrorMessage(error)) + logIntegrationEvent('remote stream token check failed', { + workspaceId, + error: toErrorMessage(error) + }) }) return { @@ -556,6 +588,28 @@ function injectionDeduplicationKey(projectId: string, event: ChangeEvent, matche function logIntegrationEvent(message: string, metadata: Record): void { console.info(`[integration-events] ${message}`, metadata) + if (isTestProcess()) return + void appendIntegrationEventLog(message, metadata) +} + +function isTestProcess(): boolean { + return process.env.NODE_ENV === 'test' || + process.env.VITEST === 'true' || + process.argv.some((arg) => arg === '--test' || arg.includes('/vitest/')) +} + +async function appendIntegrationEventLog(message: string, metadata: Record): Promise { + const entry = { + timestamp: new Date().toISOString(), + message, + metadata + } + try { + await mkdir(dirname(INTEGRATION_EVENT_LOG_PATH), { recursive: true }) + await appendFile(INTEGRATION_EVENT_LOG_PATH, `${JSON.stringify(entry)}\n`, 'utf8') + } catch { + // Diagnostics must never affect event delivery. + } } export function integrationSubscriptionSummaries( @@ -592,6 +646,8 @@ function shouldNotifyRelayfilePath(pathValue: string): boolean { leaf === '_index.json' || leaf === '.schema.json' || leaf === '.create.example.json' || + path === '/discovery' || + path.startsWith('/discovery/') || path.includes('/discovery/') || path.includes('/.relay/') || path.includes('/.relayfile-') || @@ -682,6 +738,7 @@ export class IntegrationEventBridge { const handle = await this.getWorkspaceHandle() const signature = JSON.stringify({ workspaceId: handle.workspaceId, + localMountWorkspaceId: handle.localMountWorkspaceId, watches, specs: specs.map((spec) => ({ integrationId: spec.integrationId, @@ -698,6 +755,7 @@ export class IntegrationEventBridge { logIntegrationEvent('subscribing', { projectId, workspaceId: handle.workspaceId, + localMountWorkspaceId: handle.localMountWorkspaceId, globs: watches.map((watch) => watch.glob), specs: specs.map((spec) => ({ integrationId: spec.integrationId, @@ -717,7 +775,7 @@ export class IntegrationEventBridge { path: event.resource.path }) void this.injectEvent(projectId, event, specs).catch((error) => { - console.warn('[integration-events] Event delivery failed:', { + logIntegrationEvent('event delivery failed', { projectId, eventId: event.id, error: toErrorMessage(error) @@ -731,7 +789,7 @@ export class IntegrationEventBridge { ) ) const localSubscription = watchLocalMounts( - handle.workspaceId, + handle.localMountWorkspaceId, subscribed, watches.map((watch) => watch.glob), (event) => { @@ -743,7 +801,7 @@ export class IntegrationEventBridge { source: 'local-mount' }) void this.injectEvent(projectId, event, specs).catch((error) => { - console.warn('[integration-events] Local event delivery failed:', { + logIntegrationEvent('local event delivery failed', { projectId, eventId: event.id, error: toErrorMessage(error) @@ -756,6 +814,7 @@ export class IntegrationEventBridge { logIntegrationEvent('watching local mounts', { projectId, workspaceId: handle.workspaceId, + localMountWorkspaceId: handle.localMountWorkspaceId, localRoots: localSubscription.localRoots }) subscriptions.push(localSubscription) @@ -813,16 +872,25 @@ export class IntegrationEventBridge { const bridge = await this.bridge() let allProjectAgents: string[] | null = null const recipients: string[] = [] + const listProjectAgents = async (): Promise => { + allProjectAgents ??= (await bridge.listAgents(projectId)) + .filter((agent) => agent.projectId === undefined || agent.projectId === projectId) + .map((agent) => agent.name) + return allProjectAgents + } for (const spec of matchedSpecs) { - const explicitTargets = dedupeStrings([...spec.targets.agents, ...spec.targets.channels]) + const projectAgents = spec.targets.agents.length > 0 + ? await listProjectAgents() + : null + const onlineExplicitAgents = projectAgents + ? spec.targets.agents.filter((agent) => projectAgents.includes(agent)) + : [] + const explicitTargets = dedupeStrings([...onlineExplicitAgents, ...spec.targets.channels]) if (explicitTargets.length === 0) { - allProjectAgents ??= (await bridge.listAgents(projectId)) - .filter((agent) => agent.projectId === undefined || agent.projectId === projectId) - .map((agent) => agent.name) - recipients.push(...allProjectAgents) + recipients.push(...await listProjectAgents()) } else { - recipients.push(...spec.targets.agents, ...spec.targets.channels) + recipients.push(...explicitTargets) } } @@ -872,12 +940,12 @@ export class IntegrationEventBridge { throw new Error('cloud-auth-required') } - const workspaceId = await getAccountWorkspaceId(accountWorkspaceReadyRetryOptions()) + const accountWorkspaceId = await getAccountWorkspaceId(accountWorkspaceReadyRetryOptions()) if ( accountIntegrationEventHandle && accountIntegrationEventHandle.apiUrl === auth.apiUrl && accountIntegrationEventHandle.accountKey === auth.accountKey && - accountIntegrationEventHandle.workspaceId === workspaceId + accountIntegrationEventHandle.accountWorkspaceId === accountWorkspaceId ) { return accountIntegrationEventHandle.handle } @@ -890,10 +958,11 @@ export class IntegrationEventBridge { cloudApiUrl: auth.apiUrl, accessToken: tokenProvider }) - const joined = await setup.joinWorkspace(workspaceId, { + const joined = await setup.joinWorkspace(accountWorkspaceId, { agentName: INTEGRATION_EVENT_AGENT_NAME, scopes: INTEGRATION_EVENT_SCOPES }) + const relayWorkspaceId = joined.workspaceId const workspaceTokenProvider = async (): Promise => { await joined.refreshToken() return joined.getToken() @@ -903,13 +972,14 @@ export class IntegrationEventBridge { token: workspaceTokenProvider }) const handle: RelayfileWorkspaceHandle = { - workspaceId, - client: () => createWorkspaceScopedEventClient(client, workspaceId, workspaceTokenProvider) + workspaceId: relayWorkspaceId, + localMountWorkspaceId: accountWorkspaceId, + client: () => createWorkspaceScopedEventClient(client, relayWorkspaceId, workspaceTokenProvider) } accountIntegrationEventHandle = { apiUrl: auth.apiUrl, accountKey: auth.accountKey, - workspaceId, + accountWorkspaceId, handle } return handle diff --git a/src/main/integration-mounts.test.ts b/src/main/integration-mounts.test.ts index 9a6db2f3..bb587d7f 100644 --- a/src/main/integration-mounts.test.ts +++ b/src/main/integration-mounts.test.ts @@ -172,6 +172,56 @@ describe('IntegrationMountManager', () => { ]) }) + it('preserves root-level discovery mounts for writeback metadata', async () => { + const manager = new IntegrationMountManager() + + await manager.ensureMounted([ + { + provider: 'slack', + mountPaths: ['/discovery/slack'] + } + ]) + + expect(mock.mountInputs).toHaveLength(1) + expect(mock.mountInputs[0]).toMatchObject({ + localDir: '/tmp/pear-home/.agentworkforce/pear/relayfile/workspaces/account-workspace-id/discovery/slack', + remotePath: '/discovery/slack', + agentName: 'pear-integrations-discovery-slack', + scopes: ['relayfile:fs:read:/discovery/slack/**', 'relayfile:fs:write:/discovery/slack/**'] + }) + expect(manager.localPathsFor('account-workspace-id', { + provider: 'slack', + mountPaths: ['/discovery/slack'] + })).toEqual([ + '/tmp/pear-home/.agentworkforce/pear/relayfile/workspaces/account-workspace-id/discovery/slack' + ]) + }) + + it('scopes bare discovery mounts to the integration provider', async () => { + const manager = new IntegrationMountManager() + + await manager.ensureMounted([ + { + provider: 'slack', + mountPaths: ['/discovery'] + } + ]) + + expect(mock.mountInputs).toHaveLength(1) + expect(mock.mountInputs[0]).toMatchObject({ + localDir: '/tmp/pear-home/.agentworkforce/pear/relayfile/workspaces/account-workspace-id/discovery/slack', + remotePath: '/discovery/slack', + agentName: 'pear-integrations-discovery-slack', + scopes: ['relayfile:fs:read:/discovery/slack/**', 'relayfile:fs:write:/discovery/slack/**'] + }) + expect(manager.localPathsFor('account-workspace-id', { + provider: 'slack', + mountPaths: ['/discovery'] + })).toEqual([ + '/tmp/pear-home/.agentworkforce/pear/relayfile/workspaces/account-workspace-id/discovery/slack' + ]) + }) + it('restarts mounted providers when the relayfile mount token reaches its refresh time', async () => { vi.useFakeTimers() vi.setSystemTime(new Date('2026-06-05T00:00:00.000Z')) @@ -214,6 +264,7 @@ describe('IntegrationMountManager', () => { }) await Promise.resolve() await Promise.resolve() + await new Promise((resolve) => setTimeout(resolve, 0)) expect(mock.mountInputs.filter((input) => input.remotePath === '/slack/channels')).toHaveLength(2) }) diff --git a/src/main/integration-mounts.ts b/src/main/integration-mounts.ts index d2468866..88d7d184 100644 --- a/src/main/integration-mounts.ts +++ b/src/main/integration-mounts.ts @@ -75,6 +75,10 @@ export function integrationProviderRoot(mountPath: string): string | null { function canonicalIntegrationMountPath(mountPath: string, provider: string): string | null { const providerSegment = sanitizePathSegment(provider.trim().toLowerCase()) const segments = remotePathSegments(mountPath) + if (segments[0] === 'discovery') { + const discoverySegments = segments.length > 1 ? segments : ['discovery', providerSegment] + return `/${discoverySegments.join('/')}` + } const withoutLegacyRoot = segments[0] === 'integrations' ? segments.slice(2) : segments.slice(1) if (segments[0] === 'integrations') { const legacyProvider = segments[1] || providerSegment diff --git a/src/main/integration-remote-paths.ts b/src/main/integration-remote-paths.ts new file mode 100644 index 00000000..d35011a3 --- /dev/null +++ b/src/main/integration-remote-paths.ts @@ -0,0 +1,40 @@ +export function normalizeRemoteDirectoryPath(remotePath: string): string | null { + const segments = remotePath.split('/').map((segment) => segment.trim()).filter(Boolean) + if (segments.some((segment) => segment === '.' || segment === '..')) return null + return `/${segments.join('/')}` +} + +export function remotePathName(remotePath: string): string { + const segments = remotePath.split('/').filter(Boolean) + return segments[segments.length - 1] || remotePath +} + +export function isRelayfilePathWithinRoot(rootPath: string, targetPath: string): boolean { + const normalizedRoot = rootPath.trim().replace(/\/+$/, '') || '/' + const normalizedTarget = targetPath.trim().replace(/\/+$/, '') || '/' + if (normalizedRoot === '/') return normalizedTarget === '/' || normalizedTarget.startsWith('/') + return normalizedTarget === normalizedRoot || normalizedTarget.startsWith(`${normalizedRoot}/`) +} + +function isImmediateNonDiscoveryParent(parentPath: string, childPath: string): boolean { + const normalizedParent = parentPath.trim().replace(/\/+$/, '') || '/' + const normalizedChild = childPath.trim().replace(/\/+$/, '') || '/' + if (normalizedChild.startsWith('/discovery/')) return false + if (normalizedParent === '/' || normalizedParent === normalizedChild) return false + if (!normalizedChild.startsWith(`${normalizedParent}/`)) return false + return normalizedChild.slice(normalizedParent.length + 1).split('/').length === 1 +} + +export function canListRemoteDirectoryForMountPaths(remotePath: string, mountPaths: string[]): boolean { + return mountPaths.some((mountPath) => + isRelayfilePathWithinRoot(mountPath, remotePath) || + isImmediateNonDiscoveryParent(remotePath, mountPath) + ) +} + +export function canShowRemoteDirectoryEntryForMountPaths(entryPath: string, mountPaths: string[]): boolean { + return mountPaths.some((mountPath) => + isRelayfilePathWithinRoot(mountPath, entryPath) || + isRelayfilePathWithinRoot(entryPath, mountPath) + ) +} diff --git a/src/main/integrations.ts b/src/main/integrations.ts index 8b50fc20..eef375a5 100644 --- a/src/main/integrations.ts +++ b/src/main/integrations.ts @@ -1,13 +1,19 @@ import { randomUUID } from 'node:crypto' import { isAbsolute, relative, resolve } from 'node:path' import { BrowserWindow, shell } from 'electron' -import type { WorkspaceHandle } from '@relayfile/sdk' +import type { TreeResponse, WorkspaceHandle } from '@relayfile/sdk' import { accountWorkspaceReadyRetryOptions, getAccountWorkspaceId, getApiUrl, resolveCloudAuth } from './auth' import { brokerManager } from './broker' import { cloudAgentManager } from './cloud-agent' import * as filesystem from './filesystem' import { integrationEventBridge, integrationSubscriptionSummaries } from './integration-event-bridge' import { integrationMountManager, integrationMountRootForWorkspace } from './integration-mounts' +import { + canListRemoteDirectoryForMountPaths, + canShowRemoteDirectoryEntryForMountPaths, + normalizeRemoteDirectoryPath, + remotePathName +} from './integration-remote-paths' import { PROJECT_INTEGRATIONS_LINK_NAME, ensureProjectIntegrationsLink, @@ -76,6 +82,11 @@ export type IntegrationOption = { hint?: string } +type SlackChannelOptionResponse = { + channels?: unknown + nextCursor?: unknown +} + export type IntegrationsEvent = | { type: 'session-update'; sessionId: string; session: IntegrationConnectSession } | { type: 'integration-added'; projectId: string; integration: ConnectedIntegration } @@ -181,7 +192,9 @@ type IntegrationSystemMessageBridge = { const POLL_INTERVAL_MS = 2_000 const POLL_TIMEOUT_MS = 5 * 60_000 const CATALOG_CACHE_MS = 5 * 60_000 +const SYSTEM_MESSAGE_DEBOUNCE_MS = 1_000 const CATALOG_PATH = '/api/v1/integrations/catalog' +const MAX_REMOTE_DIRECTORY_ENTRIES = 5_000 // Only providers currently active in ../cloud are surfaced. This mirrors the // non-deprecated relayfile providers in @@ -266,6 +279,11 @@ function isHttpStatus(error: unknown, status: number): boolean { // not exist on the server, so every mount mirrored an empty tree. Rewrite // both forms to the real root-level layout. function rewriteIntegrationMountPath(mountPath: string, relayfileProvider: string): string { + const discovery = mountPath.match(/^\/discovery(?:\/.*)?$/u) + if (discovery) { + const normalized = mountPath.replace(/\/+$/u, '') + return normalized === '/discovery' ? `/discovery/${relayfileProvider}` : normalized + } const prefixed = mountPath.match(/^\/integrations\/[^/]+(\/.*)?$/) if (prefixed) return `/${relayfileProvider}${prefixed[1] ?? ''}` const rootLevel = mountPath.match(/^\/[^/]+(\/.*)?$/) @@ -279,6 +297,10 @@ function projectIntegrationPathForRelayfilePath(mountPath: string): string { return `${PROJECT_INTEGRATIONS_LINK_NAME}${normalized}` } +function discoveryMountPathForProvider(provider: string): string { + return `/discovery/${toRelayfileProvider(provider)}` +} + function normalizeCapabilities(value: unknown): IntegrationCapabilities { const record = isRecord(value) ? value : {} return { @@ -595,6 +617,7 @@ export class IntegrationsManager { private sessions = new Map() private sessionMetadata = new Map() private pollTimers = new Map() + private systemMessageTimers = new Map() private catalogCache: IntegrationAdapter[] | null = null private catalogFetchedAt = 0 @@ -679,23 +702,66 @@ export class IntegrationsManager { return filesystem.readTextPreview(resolvedPath) } + async listRemoteDirectory(projectId: string, remotePath: string): Promise { + if (!this.findProject(projectId)) throw new Error(`Project not found: ${projectId}`) + const path = normalizeRemoteDirectoryPath(remotePath) + if (!path || path === '/') throw new Error('Integration remote directory path is required') + const mountPaths = this.listableRemoteMountPaths(projectId) + if (!canListRemoteDirectoryForMountPaths(path, mountPaths)) { + throw new Error('Integration remote directory is outside this project integration scope') + } + + return this.withWorkspaceHandle(async (handle) => { + const entries: filesystem.ExplorerEntry[] = [] + let cursor: string | undefined + + do { + const tree = await handle.client().listTree(handle.workspaceId, { + path, + depth: 1, + ...(cursor ? { cursor } : {}) + }) as TreeResponse + + for (const entry of tree.entries) { + if (entry.path === path) continue + if (!canShowRemoteDirectoryEntryForMountPaths(entry.path, mountPaths)) continue + entries.push({ + name: remotePathName(entry.path), + path: entry.path, + type: entry.type === 'dir' ? 'directory' : 'file' + }) + } + + cursor = tree.nextCursor ?? undefined + } while (cursor && entries.length < MAX_REMOTE_DIRECTORY_ENTRIES) + + return entries.sort((a, b) => { + if (a.type !== b.type) return a.type === 'directory' ? -1 : 1 + return a.name.localeCompare(b.name) + }) + }) + } + async listOptions(projectId: string, provider: string, resource: string): Promise { if (!this.findProject(projectId)) throw new Error(`Project not found: ${projectId}`) - const workspaceId = await getAccountWorkspaceId(accountWorkspaceReadyRetryOptions()) const normalizedProvider = toRelayfileProvider(provider) const normalizedResource = resource.trim().toLowerCase() if (!normalizedResource) throw new Error('Integration option resource is required') - let payload: unknown - try { - payload = await this.fetchJson( - 'GET', - `api/v1/workspaces/${workspaceId}/integrations/${encodeURIComponent(normalizedProvider)}/options/${encodeURIComponent(normalizedResource)}` - ) - } catch (error) { - if (isHttpStatus(error, 404) || /\b404\b/u.test(toErrorMessage(error))) return [] - throw error - } + const payload = await this.withWorkspaceHandle(async (handle) => { + try { + return await handle.requestJson({ + operation: 'listIntegrationOptions', + method: 'GET', + path: `api/v1/workspaces/${handle.workspaceId}/integrations/${encodeURIComponent(normalizedProvider)}/options/${encodeURIComponent(normalizedResource)}` + }) as unknown + } catch (error) { + if (normalizedProvider === 'slack' && normalizedResource === 'channels') { + return this.listLegacySlackChannelOptions(handle, normalizedProvider, error) + } + throw error + } + }) const rawOptions = isRecord(payload) && Array.isArray(payload.options) ? payload.options : [] return rawOptions @@ -709,12 +775,58 @@ export class IntegrationsManager { .filter((entry): entry is IntegrationOption => entry !== null) } + private async listLegacySlackChannelOptions( + handle: WorkspaceHandle, + provider: string, + originalError: unknown + ): Promise<{ options: IntegrationOption[] }> { + const options: IntegrationOption[] = [] + let cursor: string | undefined + + try { + do { + const query = cursor ? `?cursor=${encodeURIComponent(cursor)}` : '' + const payload = await handle.requestJson({ + operation: 'listSlackAvailableChannels', + method: 'GET', + path: `api/v1/workspaces/${handle.workspaceId}/integrations/${encodeURIComponent(provider)}/channels/available${query}` + }) as SlackChannelOptionResponse + const channels = Array.isArray(payload.channels) ? payload.channels : [] + for (const channel of channels) { + if (!isRecord(channel)) continue + const value = readString(channel.id) + if (!value) continue + const name = readString(channel.name) || value + const isPrivate = channel.isPrivate === true || channel.is_private === true + options.push({ + value, + label: `#${name}`, + ...(isPrivate ? { hint: 'private' } : {}) + }) + } + cursor = readString(payload.nextCursor) + } while (cursor) + return { options } + } catch (fallbackError) { + throw new Error( + `Slack channel options are unavailable: ${toErrorMessage(originalError)}; fallback failed: ${toErrorMessage(fallbackError)}` + ) + } + } + async startLocalMountDaemon(): Promise { await this.syncLocalMounts() await this.syncAllEventSubscriptions() } + async notifyAgentState(projectId: string): Promise { + await this.syncAgentState(projectId, true) + } + async shutdownLocalMounts(): Promise { + for (const timer of this.systemMessageTimers.values()) clearTimeout(timer) + this.systemMessageTimers.clear() + // Unlink the per-project `.integrations` symlinks before stopping the // mounts so a closed app leaves no dangling links in project trees. await Promise.all( @@ -1379,8 +1491,7 @@ export class IntegrationsManager { return dedupeStrings( this.listConnected(projectId) .filter((integration) => this.isVisibleInProject(projectId, integration.integrationId)) - .filter((integration) => integration.downloadHistoricalData === true) - .flatMap((integration) => this.canonicalMountPathsForIntegration(integration)) + .flatMap((integration) => this.mountPathsForAgentWorkspace(integration)) ) } @@ -1397,7 +1508,7 @@ export class IntegrationsManager { const subscriptionsReady = await this.syncEventSubscriptions(projectId) if (notifyAgent) { - await this.safeInjectSystemMessage(projectId, this.buildSystemMessageSnippet(integrations, subscriptionsReady)) + this.scheduleSystemMessage(projectId, this.buildSystemMessageSnippet(integrations, subscriptionsReady)) } } @@ -1415,13 +1526,16 @@ export class IntegrationsManager { const scopeSummary = scopeLabels.length > 0 ? scopeLabels.join(', ') : 'all configured scope' const mountPaths = this.canonicalMountPathsForIntegration(integration) .map(projectIntegrationPathForRelayfilePath) + const discoveryPath = projectIntegrationPathForRelayfilePath(discoveryMountPathForProvider(integration.provider)) const scopeClause = mountPaths.length > 0 ? ` (event scope ${mountPaths.join(', ')})` : ' (no event scope configured)' const historyClause = integration.downloadHistoricalData === true - ? ' Historical file download is enabled.' - : ' Historical file download is off.' - lines.push(`- ${integration.provider}: ${scopeSummary}${scopeClause}.${historyClause}`) + ? ` Historical provider records are available at ${mountPaths.join(', ') || 'the configured provider paths'}.` + : ' Historical provider records are not mounted.' + lines.push( + `- ${integration.provider}: ${scopeSummary}${scopeClause}. Writeback discovery is available at ${discoveryPath}. ${historyClause}` + ) } } @@ -1429,7 +1543,7 @@ export class IntegrationsManager { lines.push('') if (!subscriptionsReady && subscriptions.length > 0) { lines.push('Integration event subscriptions are requested for this project, but Pear could not register them with Relayfile yet.') - lines.push('Do not assume notifications will arrive until a later integrations update confirms active subscriptions; read the mounted integration files when the user asks for current state.') + lines.push('Do not assume notifications will arrive until a later integrations update confirms active subscriptions. If historical provider records are mounted for an integration, read them when the user asks for current state; otherwise rely on incoming events and explicit user-provided context.') } else if (subscriptions.length === 0) { lines.push('No integration event subscriptions are active for this project.') } else { @@ -1447,12 +1561,34 @@ export class IntegrationsManager { } lines.push( - `Historical files are downloaded and mounted through ${PROJECT_INTEGRATIONS_LINK_NAME}/ only for integrations where historical download is enabled. Incoming webhook events do not require downloading history.`, + `Writeback discovery schemas and examples are mounted through ${PROJECT_INTEGRATIONS_LINK_NAME}/discovery// for connected integrations. Historical provider records are mounted through ${PROJECT_INTEGRATIONS_LINK_NAME}/ only when historical download is enabled. Incoming webhook events do not require downloading history.`, '' ) return lines.join('\n') } + private mountPathsForAgentWorkspace(integration: ConnectedIntegration): string[] { + return dedupeStrings([ + discoveryMountPathForProvider(integration.provider), + ...(integration.downloadHistoricalData === true + ? this.canonicalMountPathsForIntegration(integration) + : []) + ]) + } + + private mountPathsForLocalSync(integration: ConnectedIntegration): string[] { + return this.mountPathsForAgentWorkspace(integration) + } + + private listableRemoteMountPaths(projectId: string): string[] { + return dedupeStrings( + this.visibleIntegrationsForProject(projectId).flatMap((integration) => [ + discoveryMountPathForProvider(integration.provider), + ...this.canonicalMountPathsForIntegration(integration) + ]) + ) + } + private canonicalMountPathsForIntegration(integration: ConnectedIntegration): string[] { return dedupeStrings( integration.mountPaths.map((mountPath) => @@ -1485,6 +1621,17 @@ export class IntegrationsManager { } } + private scheduleSystemMessage(projectId: string, message: string): void { + const existing = this.systemMessageTimers.get(projectId) + if (existing) clearTimeout(existing) + + const timer = setTimeout(() => { + this.systemMessageTimers.delete(projectId) + void this.safeInjectSystemMessage(projectId, message) + }, SYSTEM_MESSAGE_DEBOUNCE_MS) + this.systemMessageTimers.set(projectId, timer) + } + private async safeUpdateMountPaths(projectId: string, paths: string[]): Promise { try { const bridge: CloudAgentIntegrationsBridge = cloudAgentManager @@ -1518,17 +1665,24 @@ export class IntegrationsManager { integration })) ) - const byProvider = new Map() + const byProvider = new Map }>() for (const { projectId, integration } of localEntries) { if (!this.isVisibleInProject(projectId, integration.integrationId)) continue - if (integration.downloadHistoricalData !== true) continue - byProvider.set(`${toRelayfileProvider(integration.provider)}:${integration.integrationId}`, integration) + const key = `${toRelayfileProvider(integration.provider)}:${integration.integrationId}` + const entry = byProvider.get(key) ?? { + integration, + mountPaths: new Set() + } + for (const mountPath of this.mountPathsForLocalSync(integration)) { + entry.mountPaths.add(mountPath) + } + byProvider.set(key, entry) } const integrations = Array.from(byProvider.values()) try { - await integrationMountManager.ensureMounted(integrations.map((integration) => ({ + await integrationMountManager.ensureMounted(integrations.map(({ integration, mountPaths }) => ({ provider: integration.provider, - mountPaths: this.canonicalMountPathsForIntegration(integration) + mountPaths: Array.from(mountPaths) }))) } catch (error) { console.warn('[integrations] Failed to reconcile local integration mount:', toErrorMessage(error)) @@ -1563,17 +1717,11 @@ export class IntegrationsManager { const workspaceId = integrationMountManager.currentWorkspaceId() if (!workspaceId) return integrations return integrations.map((integration) => { - if (integration.downloadHistoricalData !== true) { - return { - ...integration, - localMountPaths: [] - } - } return { ...integration, localMountPaths: integrationMountManager.localPathsFor(workspaceId, { provider: integration.provider, - mountPaths: this.canonicalMountPathsForIntegration(integration) + mountPaths: this.mountPathsForLocalSync(integration) }) } }) diff --git a/src/main/ipc-handlers.ts b/src/main/ipc-handlers.ts index 4ce14e5d..16d0cf2b 100644 --- a/src/main/ipc-handlers.ts +++ b/src/main/ipc-handlers.ts @@ -179,6 +179,9 @@ export function registerIpcHandlers(): void { return false } await brokerManager.start(projectId, cwd, name, win, channels) + void integrationsManager.notifyAgentState(projectId).catch((error) => { + console.warn('[integrations] Failed to notify agents after broker start:', error instanceof Error ? error.message : String(error)) + }) return true }) @@ -560,6 +563,10 @@ export function registerIpcHandlers(): void { return integrationsManager.listMountDirectory(projectId, integrationId, dirPath) }) + ipcMain.handle('integrations:list-remote-dir', async (_, projectId: string, remotePath: string) => { + return integrationsManager.listRemoteDirectory(projectId, remotePath) + }) + ipcMain.handle('integrations:read-mount-preview', async (_, projectId: string, integrationId: string, filePath: string) => { return integrationsManager.readMountPreview(projectId, integrationId, filePath) }) diff --git a/src/preload/index.ts b/src/preload/index.ts index 945ebbb4..1628ce1c 100644 --- a/src/preload/index.ts +++ b/src/preload/index.ts @@ -363,6 +363,8 @@ const api = { list: (projectId: string) => invoke('integrations:list', projectId), listMountDir: (projectId: string, integrationId: string, dirPath: string) => invoke('integrations:list-mount-dir', projectId, integrationId, dirPath), + listRemoteDir: (projectId: string, remotePath: string) => + invoke('integrations:list-remote-dir', projectId, remotePath), readMountPreview: (projectId: string, integrationId: string, filePath: string) => invoke('integrations:read-mount-preview', projectId, integrationId, filePath), listOptions: (projectId: string, provider: string, resource: string) => diff --git a/src/renderer/src/components/settings/AccountSettings.tsx b/src/renderer/src/components/settings/AccountSettings.tsx index c90e0612..a029a156 100644 --- a/src/renderer/src/components/settings/AccountSettings.tsx +++ b/src/renderer/src/components/settings/AccountSettings.tsx @@ -139,8 +139,21 @@ function relativeMountPath(root: string | null, path: string | null): string { : normalizedPath } +function canBrowseSyncedData(integration: ConnectedIntegration): boolean { + return integration.downloadHistoricalData === true +} + +function localPathSegments(path: string): string[] { + return path.split(/[\\/]+/u).filter(Boolean) +} + +function historicalLocalMountPaths(integration: ConnectedIntegration): string[] { + return (integration.localMountPaths || []) + .filter((root) => !localPathSegments(root).includes('discovery')) +} + function IntegrationRelayfileBrowser({ - integration, + roots, currentPath, entries, preview, @@ -150,7 +163,7 @@ function IntegrationRelayfileBrowser({ onOpenFile, onRefresh }: { - integration: ConnectedIntegration + roots: string[] currentPath: string | null entries: FsDirEntry[] preview: { path: string; result: FsReadPreviewResult } | null @@ -160,7 +173,6 @@ function IntegrationRelayfileBrowser({ onOpenFile: (path: string) => void onRefresh: () => void }): React.ReactNode { - const roots = integration.localMountPaths || [] const currentRoot = rootForPath(roots, currentPath) const parent = currentPath ? pathParent(currentPath) : null const canGoUp = !!parent && !!currentRoot && pathWithinRoot(currentRoot, parent) @@ -395,6 +407,18 @@ export function AccountSettings(): React.ReactNode { }) }, [activeProjectId, completeSession, loadConnected]) + useEffect(() => { + if (!expandedIntegrationId) return + const expanded = connected.find((integration) => integration.integrationId === expandedIntegrationId) + if (expanded && canBrowseSyncedData(expanded)) return + + setExpandedIntegrationId(null) + setBrowserPath(null) + setBrowserEntries([]) + setBrowserPreview(null) + setBrowserError(null) + }, [connected, expandedIntegrationId]) + const startConnect = useCallback(async (adapter: IntegrationAdapter) => { if (!activeProjectId) { setError('Select a project before connecting an integration.') @@ -501,6 +525,9 @@ export function AccountSettings(): React.ReactNode { }, [activeProjectId]) const toggleMountBrowser = useCallback((integration: ConnectedIntegration) => { + if (!canBrowseSyncedData(integration)) return + + const roots = historicalLocalMountPaths(integration) if (expandedIntegrationId === integration.integrationId) { setExpandedIntegrationId(null) setBrowserPath(null) @@ -510,7 +537,7 @@ export function AccountSettings(): React.ReactNode { return } - const root = integration.localMountPaths?.[0] + const root = roots[0] setExpandedIntegrationId(integration.integrationId) setBrowserPath(root || null) setBrowserEntries([]) @@ -599,10 +626,15 @@ export function AccountSettings(): React.ReactNode {

Connected

- {connected.length} + {loading ? 'Loading' : connected.length}
- {connected.length === 0 ? ( + {loading && connected.length === 0 ? ( +
+ + Loading integrations +
+ ) : connected.length === 0 ? (
No integrations connected
@@ -610,34 +642,50 @@ export function AccountSettings(): React.ReactNode { connected.map((integration) => { const adapter = adapterByProvider.get(canonicalProviderKey(integration.provider)) const expanded = expandedIntegrationId === integration.integrationId + const browsable = canBrowseSyncedData(integration) + const browserRoots = historicalLocalMountPaths(integration) return (
- + ) : ( +
+ +
+
+ {adapter?.displayName || integration.provider} +
- )} +
- + )}
- {expanded && ( + {expanded && browsable && ( void openMountDirectory(integration, path)} onOpenFile={(path) => void openMountPreview(integration, path)} onRefresh={() => { - const root = integration.localMountPaths?.[0] || null + const root = browserRoots[0] || null if (browserPath) { void openMountDirectory(integration, browserPath) } else if (root) { diff --git a/src/renderer/src/components/settings/ProjectSettings.tsx b/src/renderer/src/components/settings/ProjectSettings.tsx index 06b775d1..047a3e4d 100644 --- a/src/renderer/src/components/settings/ProjectSettings.tsx +++ b/src/renderer/src/components/settings/ProjectSettings.tsx @@ -1,5 +1,5 @@ import type React from 'react' -import { useCallback, useEffect, useMemo, useState } from 'react' +import { useCallback, useEffect, useMemo, useRef, useState } from 'react' import { AlertTriangle, Bell, @@ -234,9 +234,7 @@ function providerLocalMountPath(integration: ConnectedIntegration, provider: str const providerSegment = provider.trim().toLowerCase() const providerRoots = (integration.localMountPaths || []) .filter((root) => localPathSegments(root).at(-1)?.toLowerCase() === providerSegment) - return providerRoots.find((root) => !localPathSegments(root).includes('discovery')) - || providerRoots[0] - || null + return providerRoots.find((root) => !localPathSegments(root).includes('discovery')) || null } function joinLocalPath(root: string, ...segments: string[]): string { @@ -247,6 +245,71 @@ function joinLocalPath(root: string, ...segments: string[]): string { return suffix ? `${root.replace(/[/\\]+$/u, '')}/${suffix}` : root } +function slackChannelResourceFromEntry(entry: { name: string; path: string }): IntegrationAccessibleResource { + const channelId = entry.name.includes('__') + ? entry.name.split('__')[0] + : entry.name.includes('--') + ? entry.name.split('--').at(-1) || entry.name + : entry.name + const name = entry.name.includes('__') + ? entry.name.split('__').slice(1).join('__') + : entry.name.includes('--') + ? entry.name.split('--').slice(0, -1).join('--') + : entry.name + const remotePath = `/slack/channels/${entry.name}` + return { + id: remotePath, + displayName: name || channelId, + name: name || channelId, + path: entry.path, + metadata: { + channelFolder: entry.name, + channelId, + remotePath + } + } +} + +function isConcreteSlackChannelEntry(entry: { name: string; path: string }): boolean { + if (!entry.name || entry.name.startsWith('_') || entry.name === 'by-name') return false + if (entry.name.includes('{') || entry.name.includes('}')) return false + if (localPathSegments(entry.path).includes('discovery')) return false + return true +} + +function slackChannelResourceFromOption(option: { value: string; label: string; hint?: string }): IntegrationAccessibleResource { + const id = option.value.trim() + const name = option.label.replace(/^#/u, '').trim() || id + return { + id, + displayName: option.label, + name, + metadata: { + channelId: id, + ...(option.hint ? { hint: option.hint } : {}) + } + } +} + +function integrationStatusSummary( + integration: ConnectedIntegration, + visibility: ProjectVisibilityMetadata, + visible: boolean +): string { + if (!visible) return 'Hidden' + if (integration.downloadHistoricalData === true) return visibilityMountPaths(integration, visibility).join(', ') + if (integration.subscribeAgent === true) return 'Listening for events' + return 'Historical files not mounted' +} + +const RESOURCE_CACHE_TTL_MS = 5 * 60 * 1000 + +type ResourceCacheEntry = { + expiresAt: number + resources?: IntegrationAccessibleResource[] + promise?: Promise +} + function notificationTargetValue(scope: Record): string { const agent = firstScopeString(scope, NOTIFICATION_AGENT_SCOPE_KEYS) if (agent) return `agent:${agent.replace(/^@/u, '')}` @@ -255,6 +318,10 @@ function notificationTargetValue(scope: Record): string { return 'all' } +function resolvedNotificationTargetValue(value: string, knownValues: Set): string { + return knownValues.has(value) ? value : 'all' +} + function clearNotificationTargetScope(scope: Record): Record { const nextScope = { ...scope } for (const key of [...NOTIFICATION_AGENT_SCOPE_KEYS, ...NOTIFICATION_CHANNEL_SCOPE_KEYS]) { @@ -278,6 +345,7 @@ function IntegrationVisibilitySection({ const [error, setError] = useState(null) const [scopeEditorIntegrationId, setScopeEditorIntegrationId] = useState(null) const [pendingScopeValue, setPendingScopeValue] = useState(null) + const resourceCacheRef = useRef(new Map()) const load = useCallback(async () => { setError(null) @@ -407,53 +475,89 @@ function IntegrationVisibilitySection({ } }, [projectId]) + const cachedResources = useCallback( + async (cacheKey: string, loadResources: () => Promise): Promise => { + const now = Date.now() + const cached = resourceCacheRef.current.get(cacheKey) + if (cached?.resources && cached.expiresAt > now) return cached.resources + if (cached?.promise) return cached.promise + + const promise = loadResources() + .then((resources) => { + resourceCacheRef.current.set(cacheKey, { + resources, + expiresAt: Date.now() + RESOURCE_CACHE_TTL_MS + }) + return resources + }) + .catch((err) => { + resourceCacheRef.current.delete(cacheKey) + throw err + }) + + resourceCacheRef.current.set(cacheKey, { + promise, + expiresAt: now + RESOURCE_CACHE_TTL_MS + }) + return promise + }, + [] + ) + const listSlackChannelResources = useCallback(async (integration: ConnectedIntegration): Promise => { - const listMountedSlackChannels = async (): Promise => { - const slackRoot = providerLocalMountPath(integration, 'slack') - if (!slackRoot) throw new Error('Slack Relayfile mount is not available yet.') - const channelsPath = joinLocalPath(slackRoot, 'channels') - const entries = await pear.integrations.listMountDir(projectId, integration.integrationId, channelsPath) - return entries - .filter((entry) => entry.type === 'directory' && !entry.name.startsWith('_') && entry.name !== 'by-name') - .map((entry) => { - const id = entry.name.includes('__') - ? entry.name.split('__')[0] - : entry.name.includes('--') - ? entry.name.split('--').at(-1) || entry.name - : entry.name - const name = entry.name.includes('__') - ? entry.name.split('__').slice(1).join('__') - : entry.name.includes('--') - ? entry.name.split('--').slice(0, -1).join('--') - : entry.name - return { - id, - displayName: name || id, - name: name || id, + const cacheKey = `slack-channels:${projectId}:${integration.integrationId}` + return cachedResources(cacheKey, async () => { + const listRemoteSlackChannels = async (): Promise => { + const listRemoteDir = (pear.integrations as typeof pear.integrations & { + listRemoteDir?: typeof pear.integrations.listRemoteDir + }).listRemoteDir + if (typeof listRemoteDir !== 'function') throw new Error('Slack Relayfile remote directory listing is not available yet.') + + const entries = await listRemoteDir(projectId, '/slack/channels') + return entries + .filter((entry) => entry.type === 'directory' && isConcreteSlackChannelEntry(entry)) + .map(slackChannelResourceFromEntry) + } + + const listMountedSlackChannels = async (): Promise => { + const slackRoot = providerLocalMountPath(integration, 'slack') + if (!slackRoot) throw new Error('Slack Relayfile mount is not available yet.') + const channelsPath = joinLocalPath(slackRoot, 'channels') + const entries = await pear.integrations.listMountDir(projectId, integration.integrationId, channelsPath) + return entries + .filter((entry) => entry.type === 'directory' && isConcreteSlackChannelEntry(entry)) + .map((entry) => slackChannelResourceFromEntry({ + name: entry.name, path: joinLocalPath(channelsPath, entry.name) - } - }) - } - const listOptions = (pear.integrations as typeof pear.integrations & { - listOptions?: typeof pear.integrations.listOptions - }).listOptions - if (typeof listOptions !== 'function') { - return listMountedSlackChannels() - } + })) + } + const listOptions = (pear.integrations as typeof pear.integrations & { + listOptions?: typeof pear.integrations.listOptions + }).listOptions - try { - const options = await listOptions(projectId, integration.provider, 'channels') - if (options.length === 0) return listMountedSlackChannels() - return options.map((option) => ({ - id: option.value, - displayName: option.label, - name: option.label.replace(/^#/u, ''), - metadata: option.hint ? { hint: option.hint } : undefined - })) - } catch { - return listMountedSlackChannels() - } - }, [projectId]) + if (typeof listOptions === 'function') { + try { + const options = await listOptions(projectId, integration.provider, 'channels') + const optionChannels = options.map(slackChannelResourceFromOption) + if (optionChannels.length > 0) return optionChannels + } catch (err) { + console.warn('[integrations] Failed to list Slack channel options:', err) + const mountedChannels = await listMountedSlackChannels().catch(() => []) + if (mountedChannels.length > 0) return mountedChannels + const message = err instanceof Error ? err.message : String(err) + throw new Error(`Slack channel options are unavailable: ${message}`) + } + } + + const remoteChannels = await listRemoteSlackChannels().catch((err) => { + console.warn('[integrations] Failed to list remote Slack channels:', err) + return [] + }) + if (remoteChannels.length > 0) return remoteChannels + + return listMountedSlackChannels().catch(() => []) + }) + }, [cachedResources, projectId]) const saveSlackSourceChannels = useCallback(async (integration: ConnectedIntegration) => { if (!pendingScopeValue) return @@ -499,7 +603,12 @@ function IntegrationVisibilitySection({
)}
- {integrations.length === 0 ? ( + {loading && integrations.length === 0 ? ( +
+ + Loading integrations +
+ ) : integrations.length === 0 ? (
No account integrations connected
@@ -510,14 +619,21 @@ function IntegrationVisibilitySection({ const subscribed = integration.subscribeAgent === true const historyDownload = integration.downloadHistoricalData === true const busy = busyIntegrationId === integration.integrationId - const notificationTarget = notificationTargetValue(integration.scope) const slack = isSlackProvider(integration.provider) const scopeEditorOpen = scopeEditorIntegrationId === integration.integrationId + const selectedSlackSourceIds = Array.from(new Set([ + ...integration.mountPaths, + ...scopeStringList(integration.scope, 'channels') + ])) const knownTargetValues = new Set([ 'all', ...agentNames.map((agent) => `agent:${agent}`), ...channels.map((channel) => `channel:${channel}`) ]) + const notificationTarget = resolvedNotificationTargetValue( + notificationTargetValue(integration.scope), + knownTargetValues + ) return (
@@ -526,11 +642,9 @@ function IntegrationVisibilitySection({
{integration.provider}
- {visible - ? visibilityMountPaths(integration, visibility).join(', ') - : providerMountRoot(integration.provider)} + {integrationStatusSummary(integration, visibility, visible)}
- {integration.localMountPaths && integration.localMountPaths.length > 0 && ( + {historyDownload && integration.localMountPaths && integration.localMountPaths.length > 0 && (
{integration.localMountPaths.join(', ')}
@@ -595,9 +709,6 @@ function IntegrationVisibilitySection({ title="Integration event delivery target" aria-label="Integration event delivery target" > - {!knownTargetValues.has(notificationTarget) && ( - - )} {agentNames.length > 0 && ( @@ -636,7 +747,7 @@ function IntegrationVisibilitySection({ listSlackChannelResources(integration)} onChange={setPendingScopeValue} /> diff --git a/src/renderer/src/components/settings/scope-pickers/SlackChannelPicker.tsx b/src/renderer/src/components/settings/scope-pickers/SlackChannelPicker.tsx index a0741be0..7a102cf3 100644 --- a/src/renderer/src/components/settings/scope-pickers/SlackChannelPicker.tsx +++ b/src/renderer/src/components/settings/scope-pickers/SlackChannelPicker.tsx @@ -19,6 +19,13 @@ function slackPathSlug(value: string): string { } function channelMountSegment(resource: Parameters[0]): string { + const channelFolder = metadataText(resource, 'channelFolder') + if (channelFolder) return channelFolder + + const path = resourceText(resource, 'path') + const pathMatch = path.match(/^\/?slack\/channels\/([^/]+)$/u) + if (pathMatch?.[1]) return pathMatch[1] + const id = resourceText(resource, 'id') const slug = slackPathSlug(channelName(resource)) if (!id) return slug @@ -39,7 +46,7 @@ export function SlackChannelPicker(props: ScopePickerProps): React.ReactNode { }} getResourceDescription={(resource) => metadataText(resource, 'workspace', 'team') || resourceText(resource, 'path')} getResourceMountSegment={channelMountSegment} - getResourceScopeId={(resource) => resourceText(resource, 'id', 'slug', 'name')} + getResourceScopeId={(resource) => metadataText(resource, 'channelId') || resourceText(resource, 'id', 'slug', 'name')} /> ) } diff --git a/src/shared/types/ipc.ts b/src/shared/types/ipc.ts index 81ed3383..889af107 100644 --- a/src/shared/types/ipc.ts +++ b/src/shared/types/ipc.ts @@ -851,6 +851,7 @@ export interface PearAPI { catalog: () => Promise list: (projectId: string) => Promise listMountDir: (projectId: string, integrationId: string, dirPath: string) => Promise + listRemoteDir: (projectId: string, remotePath: string) => Promise readMountPreview: (projectId: string, integrationId: string, filePath: string) => Promise listOptions: (projectId: string, provider: string, resource: string) => Promise startConnect: (projectId: string, provider: string) => Promise