From 581068362fdf551caadcdfead06d91eedefed322 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 14 Feb 2023 13:44:03 +0000 Subject: [PATCH] The integration service now returns the response in the format the webapp is expecting --- apps/integrations/src/api/v1/action/index.ts | 115 ++++++++++++++---- apps/webapp/app/env.server.ts | 2 + .../performIntegrationRequest.server.ts | 113 +++++++++++++---- 3 files changed, 183 insertions(+), 47 deletions(-) diff --git a/apps/integrations/src/api/v1/action/index.ts b/apps/integrations/src/api/v1/action/index.ts index 4270b71a0..043f6e80e 100644 --- a/apps/integrations/src/api/v1/action/index.ts +++ b/apps/integrations/src/api/v1/action/index.ts @@ -1,7 +1,10 @@ import { PostgresCacheService } from "cache/postgresCache"; import { AuthCredentials } from "core/authentication/types"; +import { RequestError } from "core/request/errors"; +import { Service } from "core/service/types"; import { validateInputs } from "core/validation/inputs"; import { Request, Response } from "express"; +import { stat } from "fs"; import { catalog } from "integrations/catalog"; import { z } from "zod"; @@ -15,6 +18,17 @@ const requestBodySchema = z.object({ }), }); +type ReturnResponse = { + response: NormalizedResponse; + isRetryable: boolean; + ok: boolean; +}; + +type NormalizedResponse = { + output: NonNullable; + context: any; +}; + export async function handleAction(req: Request, res: Response) { const { service, action } = req.params; @@ -24,11 +38,13 @@ export async function handleAction(req: Request, res: Response) { if (!matchingService) { res.status(404).send( - JSON.stringify({ - success: false, - service, - error: { type: "missing_service", message: "Service not found" }, - }) + JSON.stringify( + error(404, false, { + type: "missing_service", + message: "Service not found", + service, + }) + ) ); return; } @@ -39,12 +55,14 @@ export async function handleAction(req: Request, res: Response) { if (!matchingAction) { res.status(404).send( - JSON.stringify({ - success: false, - service, - action, - error: { type: "missing_action", message: "Action not found" }, - }) + JSON.stringify( + error(404, false, { + type: "missing_action", + message: "Action not found", + service, + action, + }) + ) ); return; } @@ -53,10 +71,15 @@ export async function handleAction(req: Request, res: Response) { if (!parsedRequestBody.success) { res.status(400).send( - JSON.stringify({ - success: false, - error: { type: "invalid_body", issues: parsedRequestBody.error.issues }, - }) + JSON.stringify( + error(400, false, { + type: "invalid_body", + message: "Action not found", + service, + action, + issues: parsedRequestBody.error.issues, + }) + ) ); return; } @@ -127,12 +150,9 @@ export async function handleAction(req: Request, res: Response) { } ); if (!inputValidationResult.success) { - res.status(400).send( - JSON.stringify({ - success: false, - error: inputValidationResult.error, - }) - ); + res + .status(400) + .send(JSON.stringify(error(400, false, inputValidationResult.error))); return; } @@ -145,11 +165,60 @@ export async function handleAction(req: Request, res: Response) { cache, metadata ); - res.send(JSON.stringify(data)); + + //convert into the format for the webapp + const response: ReturnResponse = { + ok: true, + isRetryable: isRetryable(matchingService, data.status), + response: { + output: data.body ?? {}, + context: { + statusCode: data.status, + headers: data.headers, + }, + }, + }; + res.send(JSON.stringify(response)); } catch (e: any) { console.error(e); + + if (e instanceof Error) { + res + .status(500) + .send(JSON.stringify(error(500, false, { error: JSON.stringify(e) }))); + return; + } + + if ("error" in e) { + res.status(500).send(JSON.stringify(error(500, false, e.error))); + return; + } + res .status(500) - .send(JSON.stringify({ success: false, errors: e.toString() })); + .send(JSON.stringify(error(500, false, { error: JSON.stringify(e) }))); } } + +function isRetryable(service: Service, status: number): boolean { + return service.retryableStatusCodes.includes(status); +} + +function error( + status: number, + isRetryable: boolean, + error: Record +): ReturnResponse { + const response: ReturnResponse = { + ok: false, + isRetryable, + response: { + output: error, + context: { + statusCode: status, + headers: {}, + }, + }, + }; + return response; +} diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 93399e5de..ba27bb4c1 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -49,6 +49,8 @@ const EnvironmentSchema = z.object({ PULSAR_AUDIENCE: z.string().optional(), PULSAR_DEBUG: z.string().optional(), INTERNAL_TRIGGER_API_KEY: z.string().optional(), + INTEGRATIONS_API_KEY: z.string(), + INTEGRATIONS_API_ORIGIN: z.string(), }); export type Environment = z.infer; diff --git a/apps/webapp/app/services/requests/performIntegrationRequest.server.ts b/apps/webapp/app/services/requests/performIntegrationRequest.server.ts index 12a9a3928..fef176b2f 100644 --- a/apps/webapp/app/services/requests/performIntegrationRequest.server.ts +++ b/apps/webapp/app/services/requests/performIntegrationRequest.server.ts @@ -10,6 +10,7 @@ import type { IntegrationRequest } from "~/models/integrationRequest.server"; import { getAccessInfo } from "../accessInfo.server"; import { RedisCacheService } from "../cacheService.server"; import { getIntegrations } from "~/models/integrations.server"; +import { env } from "~/env.server"; type CallResponse = | { @@ -65,7 +66,8 @@ export class PerformIntegrationRequest { accessInfo, integrationRequest, cache, - integrationRequest.externalService.workflowId + integrationRequest.externalService.workflowId, + integrationRequest.externalService.connection.id ); if (performedRequest.ok) { @@ -224,31 +226,94 @@ export class PerformIntegrationRequest { accessInfo: AccessInfo, integrationRequest: IntegrationRequest, cache: CacheService, - workflowId: string + workflowId: string, + connectionId: string ): Promise { - const integrationInfo = getIntegrations(true).find( - (i) => i.metadata.slug === service - ); + switch (integrationRequest.version) { + case "1": { + const integrationInfo = getIntegrations(true).find( + (i) => i.metadata.slug === service + ); - if (!integrationInfo) { - throw new Error(`Unknown service: ${service}`); + if (!integrationInfo) { + throw new Error(`Unknown service: ${service}`); + } + + const { requests } = integrationInfo; + + if (!requests) { + throw new Error(`Service ${service} does not support requests`); + } + + return requests.perform({ + accessInfo, + endpoint: integrationRequest.endpoint, + params: integrationRequest.params, + cache, + metadata: { + requestId: integrationRequest.id, + workflowId: workflowId, + }, + }); + } + case "2": { + let credentials: { accessToken: string } | undefined = undefined; + switch (accessInfo.type) { + case "oauth2": + credentials = { + accessToken: accessInfo.accessToken, + }; + break; + case "api_key": + credentials = { + accessToken: accessInfo.api_key, + }; + } + + try { + const response = await fetch( + `${env.INTEGRATIONS_API_ORIGIN}/api/v1/${service}/action/${integrationRequest.endpoint}`, + { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Bearer ${env.INTEGRATIONS_API_KEY}`, + }, + body: JSON.stringify({ + credentials, + params: integrationRequest.params, + metadata: { + requestId: integrationRequest.id, + workflowId: workflowId, + connectionId, + }, + }), + } + ); + + const json = await response.json(); + return json; + } catch (e) { + console.error(e); + return { + ok: false, + isRetryable: true, + response: { + output: { + error: { + message: JSON.stringify(e), + }, + }, + context: {}, + }, + }; + } + } + default: { + throw new Error( + `Unknown integration request version: ${integrationRequest.version}` + ); + } } - - const { requests } = integrationInfo; - - if (!requests) { - throw new Error(`Service ${service} does not support requests`); - } - - return requests.perform({ - accessInfo, - endpoint: integrationRequest.endpoint, - params: integrationRequest.params, - cache, - metadata: { - requestId: integrationRequest.id, - workflowId: workflowId, - }, - }); } }