Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions components/google_drive/common/constants.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,30 @@ const FILES_MAX_PAGE_SIZE = 1000;
*/
const DEFAULT_SEARCH_FILES_LIMIT = 100;

const RETRYABLE_STATUS_CODES = [
422,
429,
500,
502,
503,
504,
];

const RATE_LIMIT_ERROR_REASONS = [
"rateLimitExceeded",
"userRateLimitExceeded",
];

const CHANGED_FILE_FIELDS = "kind,id,name,mimeType,parents,createdTime,modifiedTime,trashed,version,size,md5Checksum,webViewLink,lastModifyingUser";

/** Google Workspace types that `stashFile` can export to PDF. */
const PDF_EXPORTABLE_MIME_TYPES = [
"application/vnd.google-apps.document",
"application/vnd.google-apps.spreadsheet",
"application/vnd.google-apps.presentation",
"application/vnd.google-apps.drawing",
];

export {
GOOGLE_DRIVE_NOTIFICATION_SYNC,
GOOGLE_DRIVE_NOTIFICATION_ADD,
Expand Down Expand Up @@ -300,4 +324,9 @@ export {
// Files
FILES_MAX_PAGE_SIZE,
DEFAULT_SEARCH_FILES_LIMIT,
// Webhooks
RETRYABLE_STATUS_CODES,
RATE_LIMIT_ERROR_REASONS,
CHANGED_FILE_FIELDS,
PDF_EXPORTABLE_MIME_TYPES,
};
27 changes: 21 additions & 6 deletions components/google_drive/google_drive.app.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ import {
GOOGLE_DRIVE_UPDATE_TYPES,
GOOGLE_DRIVE_UPLOAD_TYPE_OPTIONS,
MY_DRIVE_VALUE,
RATE_LIMIT_ERROR_REASONS,
RETRYABLE_STATUS_CODES,
WEBHOOK_SUBSCRIPTION_EXPIRATION_TIME_MILLISECONDS,
} from "./common/constants.mjs";

Expand Down Expand Up @@ -445,15 +447,20 @@ export default {
* returned
* @param {number} [pageSize=1000] - the maximum number of changes to return
* per page
* @param {string} [fileFields] - the file fields to return for each change,
* e.g. `id,name,parents`. Defaults to the API's minimal file fields
* @yields
* @type {ChangesPage}
*/
async *listChanges(pageToken, driveId, pageSize = 1000) {
async *listChanges(pageToken, driveId, pageSize = 1000, fileFields) {
const drive = this.drive();
let changeRequest = {
pageToken,
pageSize,
};
if (fileFields) {
changeRequest.fields = `nextPageToken,newStartPageToken,changes(file(${fileFields}))`;
}

// As with many of the methods for Google Drive, we must
// pass a request of a different shape when we're requesting
Expand All @@ -468,7 +475,9 @@ export default {
}

while (true) {
const { data } = await drive.changes.list(changeRequest);
const { data } = await this.retryWithExponentialBackoff(
() => drive.changes.list(changeRequest),
);
const {
changes = [],
newStartPageToken,
Expand Down Expand Up @@ -1693,22 +1702,28 @@ export default {
const drive = this.drive();
return (await drive.accessproposals.resolve(opts)).data;
},
isRetryableError(error, statusCode) {
if (RETRYABLE_STATUS_CODES.includes(statusCode)) {
return true;
}
const errors = error.errors ?? error.response?.data?.error?.errors ?? [];
return statusCode === 403
&& errors.some(({ reason }) => RATE_LIMIT_ERROR_REASONS.includes(reason));
},
retryWithExponentialBackoff(func, maxAttempts = 3, baseDelayS = 2) {
let attempt = 0;

const execute = async () => {
try {
return await func();
} catch (error) {
// retry for error status 422
const statusCode = error.status || error.response?.status;
if (attempt >= maxAttempts || statusCode !== 422) {
if (attempt >= maxAttempts || !this.isRetryableError(error, statusCode)) {
throw error;
}

// display error message for 422 status
const errorMessage = error.message || error.response?.data?.message || error.response?.statusText || "Unknown error";
console.log(`Received 422 error: ${errorMessage}. Retrying attempt ${attempt + 1}/${maxAttempts}...`);
console.log(`Received ${statusCode} error: ${errorMessage}. Retrying attempt ${attempt + 1}/${maxAttempts}...`);

const delayMs = Math.pow(baseDelayS, attempt) * 1000;
await new Promise((resolve) => setTimeout(resolve, delayMs));
Expand Down
2 changes: 1 addition & 1 deletion components/google_drive/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@pipedream/google_drive",
"version": "3.0.1",
"version": "4.0.0",
"description": "Pipedream Google_drive Components",
"main": "google_drive.app.mjs",
"keywords": [
Expand Down
35 changes: 23 additions & 12 deletions components/google_drive/sources/common-dedupe-changes.mjs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
export default {
props: {
intervalAlert: {

Check warning on line 3 in components/google_drive/sources/common-dedupe-changes.mjs

View workflow job for this annotation

GitHub Actions / Lint Code Base

Component prop intervalAlert must have a description. See https://pipedream.com/docs/components/guidelines/#props

Check warning on line 3 in components/google_drive/sources/common-dedupe-changes.mjs

View workflow job for this annotation

GitHub Actions / Lint Code Base

Component prop intervalAlert must have a label. See https://pipedream.com/docs/components/guidelines/#props
type: "alert",
alertType: "info",
content: `This source can emit many events in quick succession while a file is being edited. By default, it will not emit another event for the same file for at least 1 minute.
Expand All @@ -24,29 +24,40 @@
_setFileIntervals(value) {
this.db.set("fileIntervals", value);
},
checkMinimumInterval(files) {
filterByMinimumInterval(files) {
const interval = this.perFileInterval;
if (!interval) return files;
if (!interval) {
return files;
}

const now = Date.now();
const minTimestamp = now - (interval * 1000 * 60);
const minTimestamp = Date.now() - (interval * 1000 * 60);
const savedData = this._getFileIntervals();
return files.filter(({ id }) => !savedData[id] || savedData[id] < minTimestamp);
},
recordFileEmits(fileIds) {
if (!this.perFileInterval || !fileIds.length) {
return;
}

const now = Date.now();
const minTimestamp = now - (this.perFileInterval * 1000 * 60);
const savedData = this._getFileIntervals();
Object.entries(savedData).forEach(([
key,
value,
]) => {
if (value < minTimestamp) delete savedData[key];
});

const filteredFiles = files.filter(({ id }) => {
const exists = !!savedData[id];
if (!exists) {
savedData[id] = now;
if (value < minTimestamp) {
delete savedData[key];
}
return !exists;
});
fileIds.forEach((id) => {
savedData[id] = now;
});
this._setFileIntervals(savedData);
},
checkMinimumInterval(files) {
const filteredFiles = this.filterByMinimumInterval(files);
this.recordFileEmits(filteredFiles.map(({ id }) => id));
return filteredFiles;
},
},
Expand Down
13 changes: 12 additions & 1 deletion components/google_drive/sources/common-webhook.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,12 @@ export default {
getUpdateTypes() {
return [];
},
/**
* File fields to request per change; `undefined` returns the API's minimal fields.
*/
getChangesFileFields() {
return undefined;
},
/**
* This method is responsible for processing a list of changed files
* according to the event source's purpose. As an abstract method, it must
Expand Down Expand Up @@ -201,7 +207,12 @@ export default {

const driveId = this.getDriveId();
const changedFilesStream =
this.googleDrive.listChanges(pageToken, driveId, this.changesPageSize);
this.googleDrive.listChanges(
pageToken,
driveId,
this.changesPageSize,
this.getChangesFileFields(),
);
for await (const changedFilesPage of changedFilesStream) {
const {
changedFiles,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,17 @@
// 2) A timer that runs on regular intervals, renewing the notification channel as needed

import {
CHANGED_FILE_FIELDS,
GOOGLE_DRIVE_MIME_TYPE_PREFIX,
GOOGLE_DRIVE_NOTIFICATION_ADD,
GOOGLE_DRIVE_NOTIFICATION_CHANGE,
GOOGLE_DRIVE_NOTIFICATION_UPDATE,
PDF_EXPORTABLE_MIME_TYPES,
} from "../../common/constants.mjs";
import commonDedupeChanges from "../common-dedupe-changes.mjs";
import common from "../common-webhook.mjs";
import { stashFile } from "../../common/utils.mjs";
import sampleEmit from "./test-event.mjs";
import md5 from "md5";

const { googleDrive } = common.props;

Expand All @@ -26,10 +28,8 @@ export default {
key: "google_drive-new-or-modified-files",
name: "New or Modified Files (Instant)",
description: "Emit new event when a file in the selected Drive is created, modified or trashed.",
version: "0.4.15",
version: "1.0.0",
type: "source",
// Dedupe events based on the "x-goog-message-number" header for the target channel:
// https://developers.google.com/drive/api/v3/push#making-watch-requests
dedupe: "unique",
props: {
...common.props,
Expand Down Expand Up @@ -66,7 +66,7 @@ export default {
includeLink: {
label: "Include Link",
type: "boolean",
description: "Upload file to your File Stash and emit temporary download link to the file. Google Workspace documents will be converted to PDF. See [the docs](https://pipedream.com/docs/connect/components/files) to learn more about working with files in Pipedream.",
description: "Upload file to your File Stash and emit temporary download link to the file. Google Workspace documents will be converted to PDF. Files that can't be downloaded emit `fileURLError` instead. See [the docs](https://pipedream.com/docs/connect/components/files) to learn more about working with files in Pipedream.",
default: false,
optional: true,
},
Expand Down Expand Up @@ -113,48 +113,61 @@ export default {
GOOGLE_DRIVE_NOTIFICATION_UPDATE,
];
},
generateMeta(data, headers) {
const {
id: fileId,
name: summary,
modifiedTime: tsString,
} = data;
const ts = Date.parse(tsString);
const eventId = headers && headers["x-goog-message-number"];

getChangesFileFields() {
return CHANGED_FILE_FIELDS;
},
generateMeta({
id, name, modifiedTime, trashed,
}) {
return {
id: md5(`${fileId}-${eventId || ts}`),
summary,
ts,
id: `${id}-${modifiedTime}-${trashed}`,
summary: name,
ts: Date.parse(modifiedTime),
};
},
async getChanges(headers) {
getChanges(headers) {
if (!headers) {
return {
change: { },
change: {},
};
}
const resourceUri = headers["x-goog-resource-uri"];
const metadata = await this.googleDrive.getFileMetadata(`${resourceUri}&fields=*`);
return {
...metadata,
change: {
state: headers["x-goog-resource-state"],
resourceURI: headers["x-goog-resource-uri"],
changed: headers["x-goog-changed"], // "Additional details about the changes. Possible values: content, parents, children, permissions"
},
};
},
async getFileLink(file) {
const { mimeType } = file;
if (mimeType.startsWith(GOOGLE_DRIVE_MIME_TYPE_PREFIX)
&& !PDF_EXPORTABLE_MIME_TYPES.includes(mimeType)) {
return {
fileURLError: `Files of type ${mimeType} can't be downloaded`,
};
}
try {
return {
fileURL: await stashFile(file, this.googleDrive, this.dir),
};
} catch (error) {
// Isolate per-file failures so one file can't block the page token; rate limits still retry
if (this.googleDrive.isRetryableError(error, error.status || error.response?.status)) {
throw error;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
console.log(`Could not upload file ${file.name} to the File Stash: ${error.message}`);
return {
fileURLError: error.message,
};
}
},
async processChanges(changedFiles, headers) {
const changes = await this.getChanges(headers);

const filteredFiles = this.checkMinimumInterval(changedFiles);
const changes = this.getChanges(headers);
const filteredFiles = this.filterByMinimumInterval(changedFiles);
const emittedFileIds = [];

for (const file of filteredFiles) {
file.parents = (await this.googleDrive.getFile(file.id, {
fields: "parents",
})).parents;

if (!this.shouldProcess(file)) {
console.log(`Skipping file ${file.name}`);
continue;
Expand All @@ -165,11 +178,13 @@ export default {
...changes,
};
if (this.includeLink) {
eventToEmit.fileURL = await stashFile(file, this.googleDrive, this.dir);
Object.assign(eventToEmit, await this.getFileLink(file));
}
const meta = this.generateMeta(file, headers);
this.$emit(eventToEmit, meta);
this.$emit(eventToEmit, this.generateMeta(file));
emittedFileIds.push(file.id);
}

this.recordFileEmits(emittedFileIds);
},
},
sampleEmit,
Expand Down
Original file line number Diff line number Diff line change
@@ -1 +1,30 @@
export default JSON.parse("{\n \"file\": {\n \"kind\": \"drive#file\",\n \"copyRequiresWriterPermission\": false,\n \"writersCanShare\": true,\n \"viewedByMe\": true,\n \"mimeType\": \"application/vnd.google-apps.spreadsheet\",\n \"exportLinks\": {\n \"application/x-vnd.oasis.opendocument.spreadsheet\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbTmw4&exportFormat=ods\",\n \"text/tab-separated-values\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaTXnViTOo&exportFormat=tsv\",\n \"application/pdf\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaBViTOo&exportFormat=pdf\",\n \"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbOo&exportFormat=xlsx\",\n \"text/csv\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbBViTOo&exportFormat=csv\",\n \"application/zip\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-VbTmw4Pt2BViTOo&exportFormat=zip\",\n \"application/vnd.oasis.opendocument.spreadsheet\": \"https://docs.google.com/spreadsheets/export?id=1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbViTOo&exportFormat=ods\"\n },\n \"parents\": [\n \"0ANs73yKKVA\"\n ],\n \"thumbnailLink\": \"https://docs.google.com/feeds/vt?gd=true&id=1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbTmw4&v=12&s=AMedNnoAAAAAZKWjrEqscucFpYCyRCJqnd0wtBiDtXYh&sz=s220\",\n \"iconLink\": \"https://drive-thirdparty.googleusercontent.com/16/type/application/vnd.google-apps.spreadsheet\",\n \"shared\": false,\n \"lastModifyingUser\": {\n \"displayName\": \"John Doe\",\n \"kind\": \"drive#user\",\n \"me\": true,\n \"permissionId\": \"077423361841532483\",\n \"emailAddress\": \"john@doe.com\",\n \"photoLink\": \"https://lh3.googleusercontent.com/a/default-user=s64\"\n },\n \"owners\": [\n {\n \"displayName\": \"John Doe\",\n \"kind\": \"drive#user\",\n \"me\": true,\n \"permissionId\": \"07742336189483\",\n \"emailAddress\": \"john@doe.com\",\n \"photoLink\": \"https://lh3.googleusercontent.com/a/default-user=s64\"\n }\n ],\n \"webViewLink\": \"https://docs.google.com/spreadsheets/d/1LmXTIQ0wqKP7T3-l9r127MaTXn3pVbTmwTOo/edit?usp=drivesdk\",\n \"size\": \"1024\",\n \"viewersCanCopyContent\": true,\n \"permissions\": [\n {\n \"id\": \"07742336184153259483\",\n \"displayName\": \"John Doe\",\n \"type\": \"user\",\n \"kind\": \"drive#permission\",\n \"photoLink\": \"https://lh3.googleusercontent.com/a/default-user=s64\",\n \"emailAddress\": \"john@doe.com\",\n \"role\": \"owner\",\n \"deleted\": false,\n \"pendingOwner\": false\n }\n ],\n \"hasThumbnail\": true,\n \"spaces\": [\n \"drive\"\n ],\n \"id\": \"1LmXTIQ0wqKP7T3-l9r127Maw4Pt2BViTOo\",\n \"name\": \"Pipedream 2213\",\n \"starred\": false,\n \"trashed\": false,\n \"explicitlyTrashed\": false,\n \"createdTime\": \"2023-02-06T15:13:33.023Z\",\n \"modifiedTime\": \"2023-06-28T12:52:54.570Z\",\n \"modifiedByMeTime\": \"2023-06-28T12:52:54.570Z\",\n \"viewedByMeTime\": \"2023-06-29T06:30:58.441Z\",\n \"quotaBytesUsed\": \"1024\",\n \"version\": \"53\",\n \"ownedByMe\": true,\n \"isAppAuthorized\": false,\n \"capabilities\": {\n \"canChangeViewersCanCopyContent\": true,\n \"canEdit\": true,\n \"canCopy\": true,\n \"canComment\": true,\n \"canAddChildren\": false,\n \"canDelete\": true,\n \"canDownload\": true,\n \"canListChildren\": false,\n \"canRemoveChildren\": false,\n \"canRename\": true,\n \"canTrash\": true,\n \"canReadRevisions\": true,\n \"canChangeCopyRequiresWriterPermission\": true,\n \"canMoveItemIntoTeamDrive\": true,\n \"canUntrash\": true,\n \"canModifyContent\": true,\n \"canMoveItemOutOfDrive\": true,\n \"canAddMyDriveParent\": false,\n \"canRemoveMyDriveParent\": true,\n \"canMoveItemWithinDrive\": true,\n \"canShare\": true,\n \"canMoveChildrenWithinDrive\": false,\n \"canModifyContentRestriction\": true,\n \"canChangeSecurityUpdateEnabled\": false,\n \"canAcceptOwnership\": false,\n \"canReadLabels\": true,\n \"canModifyLabels\": true\n },\n \"thumbnailVersion\": \"12\",\n \"modifiedByMe\": true,\n \"permissionIds\": [\n \"07742336184153259483\"\n ],\n \"linkShareMetadata\": {\n \"securityUpdateEligible\": false,\n \"securityUpdateEnabled\": true\n }\n }\n}");
export default {
"file": {
"kind": "drive#file",
"id": "1eWLQU_76tvfXptBGkdnjmPbpwhrOtkYS",
"name": "Quarterly Report.csv",
"mimeType": "text/csv",
"parents": [
"1jPHKePO9_B4jDvpBFyAs38Ji_tCDdjxI",
],
"createdTime": "2026-09-28T17:51:38.475Z",
"modifiedTime": "2026-09-28T18:04:12.918Z",
"trashed": false,
"version": "4",
"size": "2048",
"md5Checksum": "5d41402abc4b2a76b9719d911017c592",
"webViewLink": "https://drive.google.com/file/d/1eWLQU_76tvfXptBGkdnjmPbpwhrOtkYS/view?usp=drivesdk",
"lastModifyingUser": {
"kind": "drive#user",
"displayName": "John Doe",
"me": true,
"permissionId": "07742336189483",
"emailAddress": "john@doe.com",
"photoLink": "https://lh3.googleusercontent.com/a/default-user=s64",
},
},
"change": {
"state": "change",
"resourceURI": "https://www.googleapis.com/drive/v3/changes?alt=json&pageToken=1794",
},
};
Loading