Create a router to handle importing from integrations cloud task

This commit is contained in:
Hongbo Wu 2023-02-27 21:51:46 +08:00
parent 9aa65f0d33
commit 56fd3d7617
3 changed files with 119 additions and 31 deletions

View file

@ -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/<uid>/<date>/<type>-<uuid>.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
}

View file

@ -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<string | null> => {
return Promise.resolve('')
return Promise.resolve(null)
}
export = async (
integration: Integration,
pages: Page[]
): Promise<boolean> => {
return Promise.resolve(true)
return Promise.resolve(false)
}
retrieve = async (
token: string,
since = 0,
count = 100,
offset = 0
): Promise<RetrievedResult> => {
return Promise.resolve({ data: [], hasMore: false })
retrieve = async (req: RetrieveRequest): Promise<RetrievedResult> => {
return Promise.resolve({ data: [] })
}
}

View file

@ -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<RetrievedResult> => {
// const syncAt = integration.syncedAt
// ? integration.syncedAt.getTime() / 1000
// : 0
offset = 0,
}: RetrieveRequest): Promise<RetrievedResult> => {
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/<uid>/<date>/<type>-<uuid>.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(),
// })
}
}