diff --git a/packages/api/src/routers/svc/integrations.ts b/packages/api/src/routers/svc/integrations.ts index 6f9ed2b8a..269eb1163 100644 --- a/packages/api/src/routers/svc/integrations.ts +++ b/packages/api/src/routers/svc/integrations.ts @@ -10,6 +10,9 @@ import { getRepository } from '../../entity/utils' import { syncWithIntegration } from '../../services/integrations' import { buildLogger } from '../../utils/logger' import { DateFilter } from '../../utils/search' +import { DateTime } from 'luxon' +import { createGCSFile } from '../../utils/uploads' +import { v4 as uuidv4 } from 'uuid' export interface Message { type?: EntityType @@ -19,6 +22,14 @@ export interface Message { articleId?: string } +interface ImportEvent { + userId: string + integrationId: string +} + +const isImportEvent = (event: any): event is ImportEvent => + 'userId' in event && 'integrationId' in event + const logger = buildLogger('app.dispatch') export function integrationsServiceRouter() { @@ -158,6 +169,94 @@ export function integrationsServiceRouter() { res.status(500).send(err) } }) + // import pages from integration + router.post('/import', async (req, res) => { + logger.info('start to import pages from integration') + const { message: msgStr, expired } = readPushSubscription(req) + + if (!msgStr) { + return res.status(400).send('Bad Request') + } + + if (expired) { + logger.info('discarding expired message') + return res.status(200).send('Expired') + } + + const data = JSON.parse(msgStr) + if (!isImportEvent(data)) { + logger.info('Invalid message') + return res.status(400).send('Bad Request') + } + + const userId = data.userId + const integration = await getRepository(Integration).findOneBy({ + user: { id: userId }, + id: data.integrationId, + enabled: true, + type: IntegrationType.Import, + }) + if (!integration) { + logger.info('No active integration found for user', { userId }) + return res.status(200).send('No integration found') + } + + const integrationService = getIntegrationService(integration.name) + // import pages from integration + logger.info('importing pages from integration', { + integrationId: integration.id, + }) + + // write the list of urls to a csv file and upload it to gcs + // path style: imports///-.csv + const dateStr = DateTime.now().toISODate() + const fileUuid = uuidv4() + const fullPath = `imports/${integration.user.id}/${dateStr}/URL_LIST-${fileUuid}.csv` + // open a write_stream to the file + const file = createGCSFile(fullPath) + const writeStream = file.createWriteStream({ + contentType: 'text/csv', + }) + + try { + let hasMore = true + let offset = 0 + let since = integration.syncedAt?.getTime() || 0 + while (hasMore) { + // get pages from integration + const retrieved = await integrationService.retrieve({ + token: integration.token, + since, + offset: offset, + }) + const retrievedData = retrieved.data + if (retrievedData.length === 0) { + break + } + // write the list of urls, state and labels to the stream + const csvData = retrievedData.map((page) => { + const { url, state, labels } = page + return [url, state, labels?.join(',')].join(',') + }) + writeStream.write(csvData.join('\n')) + + hasMore = !!retrieved.hasMore + offset += retrievedData.length + since = retrieved.since || Date.now() + } + // update the integration's syncedAt + await getRepository(Integration).update(integration.id, { + syncedAt: since, + }) + + res.status(200).send('OK') + } catch (err) { + logger.error('import pages from integration failed', err) + res.status(500).send(err) + } finally { + writeStream.end() + } + }) return router } diff --git a/packages/api/src/services/integrations/integration.ts b/packages/api/src/services/integrations/integration.ts index 2f8ccc0db..d91e63030 100644 --- a/packages/api/src/services/integrations/integration.ts +++ b/packages/api/src/services/integrations/integration.ts @@ -9,27 +9,30 @@ export interface RetrievedData { } export interface RetrievedResult { data: RetrievedData[] - hasMore: boolean + hasMore?: boolean + since?: number +} + +export interface RetrieveRequest { + token: string + since?: number + count?: number + offset?: number } export abstract class IntegrationService { abstract name: string accessToken = async (token: string): Promise => { - return Promise.resolve('') + return Promise.resolve(null) } export = async ( integration: Integration, pages: Page[] ): Promise => { - return Promise.resolve(true) + return Promise.resolve(false) } - retrieve = async ( - token: string, - since = 0, - count = 100, - offset = 0 - ): Promise => { - return Promise.resolve({ data: [], hasMore: false }) + retrieve = async (req: RetrieveRequest): Promise => { + return Promise.resolve({ data: [] }) } } diff --git a/packages/api/src/services/integrations/pocket.ts b/packages/api/src/services/integrations/pocket.ts index 5bea482d6..f16b4fd48 100644 --- a/packages/api/src/services/integrations/pocket.ts +++ b/packages/api/src/services/integrations/pocket.ts @@ -2,6 +2,7 @@ import { IntegrationService, RetrievedDataState, RetrievedResult, + RetrieveRequest, } from './integration' import axios from 'axios' import { env } from '../../env' @@ -109,18 +110,15 @@ export class PocketIntegration extends IntegrationService { } } - retrieve = async ( - token: string, + retrieve = async ({ + token, since = 0, count = 100, - offset = 0 - ): Promise => { - // const syncAt = integration.syncedAt - // ? integration.syncedAt.getTime() / 1000 - // : 0 + offset = 0, + }: RetrieveRequest): Promise => { const pocketData = await this.retrievePocketData( token, - since, + since / 1000, count, offset ) @@ -138,19 +136,7 @@ export class PocketIntegration extends IntegrationService { return { data, hasMore: pocketData.complete !== 1, + since: pocketData.since * 1000, } - // // write the list of urls to a csv file and upload it to gcs - // // path style: imports///-.csv - // const dateStr = DateTime.now().toISODate() - // const fileUuid = uuidv4() - // const fullPath = `imports/${integration.user.id}/${dateStr}/URL_LIST-${fileUuid}.csv` - // const data = pocketItems.map((item) => item.given_url).join('\n') - // await uploadToBucket(fullPath, Buffer.from(data, 'utf-8'), { - // contentType: 'text/csv', - // }) - // // update the integration's syncedAt - // await getRepository(Integration).update(integration.id, { - // syncedAt: new Date(), - // }) } }