/home/techb158/cosmic-risk.abdallabala.com/src/services
NameSizeModeActions
audit-service.js6590644editdlrm
auth-service.js47790644editdlrm
dashboard-service.js11100644editdlrm
gate-service.js45120644editdlrm
import-service.js61950644editdlrm
integration-service.js166530644editdlrm
project-service.js59730644editdlrm
reporting-service.js86440644editdlrm
risk-service.js125050644editdlrm
sso-service.js31670644editdlrm
Edit: /home/techb158/cosmic-risk.abdallabala.com/src/services/integration-service.js (16653B)
const { prisma } = require("../lib/prisma"); const { encryptJson, decryptJson, redactToken } = require("../lib/crypto"); const { riskStatusToExternalStatus } = require("../domain/integrations"); const { scoreRisk } = require("../domain/risk-engine"); const { createPmClient } = require("../integrations/pm-clients"); const { writeAuditEvent } = require("./audit-service"); const SECRET_KEYS = new Set(["token", "apiToken", "accessToken", "refreshToken", "clientSecret", "apiKey"]); function liveModeEnabled(integration) { return integration.liveEnabled === true || process.env.COSMIC_LIVE_PM_ENABLED === "true"; } function normalizeCredentialPayload(payload = {}) { const credentials = Object.assign({}, payload.credentials || payload); delete credentials.credentials; return credentials; } function redactedCredentialSummary(credentials = {}) { return Object.fromEntries(Object.entries(credentials).map(([key, value]) => [ key, SECRET_KEYS.has(key) ? redactToken(value) : value ])); } function integrationClientView(integration) { return { id: integration.id, provider: integration.provider, workspaceId: integration.workspaceId, workspaceName: integration.workspaceName, externalProjectKey: integration.externalProjectKey, baseUrl: integration.baseUrl, authMode: integration.authMode, syncDirection: integration.syncDirection, liveEnabled: integration.liveEnabled, liveConfig: integration.liveConfig || {} }; } function scoreRiskForSync(risk) { const mitigations = risk.mitigations || []; const average = values => values.length ? Math.round(values.reduce((sum, value) => sum + Number(value || 0), 0) / values.length) : 0; return scoreRisk(Object.assign({}, risk, { mitigationProgress: average(mitigations.map(item => item.progressPercent)), mitigationEffectiveness: average(mitigations.map(item => item.effectivenessPercent)) })); } async function getIntegration(integrationId) { const integration = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: true, mappings: true, syncRuns: { orderBy: { createdAt: "desc" }, take: 10 } } }); if (!integration) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } return integration; } async function saveIntegrationCredentials({ integrationId, credentials, actor, request }) { const integration = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: true } }); if (!integration) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } const normalized = normalizeCredentialPayload(credentials); if (!Object.keys(normalized).length) return null; const existing = await prisma.oAuthToken.findFirst({ where: { organizationId: integration.workspace.organizationId, provider: integration.provider, integrationId } }); const encryptedToken = encryptJson(normalized); const redactedToken = JSON.stringify(redactedCredentialSummary(normalized)); const data = { organizationId: integration.workspace.organizationId, provider: integration.provider, integrationId, encryptedToken, redactedToken, expiresAt: normalized.expiresAt ? new Date(normalized.expiresAt) : null }; const token = existing ? await prisma.oAuthToken.update({ where: { id: existing.id }, data }) : await prisma.oAuthToken.create({ data }); await writeAuditEvent({ organizationId: integration.workspace.organizationId, workspaceId: integration.workspaceId, actorUserId: actor?.id || null, entityType: "OAuthToken", entityId: token.id, action: existing ? "update" : "create", afterJson: { provider: integration.provider, integrationId, redactedToken }, request }).catch(() => {}); return { id: token.id, provider: token.provider, redactedToken: token.redactedToken }; } async function loadIntegrationCredentials(integration) { const token = await prisma.oAuthToken.findFirst({ where: { organizationId: integration.workspace.organizationId, provider: integration.provider, integrationId: integration.id }, orderBy: { updatedAt: "desc" } }); const stored = token ? decryptJson(token.encryptedToken) : {}; return Object.assign({}, stored, integration.liveConfig || {}, { baseUrl: stored.baseUrl || integration.baseUrl || process.env.JIRA_BASE_URL, projectKey: stored.projectKey || integration.externalProjectKey, listId: stored.listId || integration.externalProjectKey, projectGid: stored.projectGid || integration.externalProjectKey, planId: stored.planId || integration.externalProjectKey, bucketId: stored.bucketId || integration.liveConfig?.bucketId }); } async function updateIntegration(integrationId, payload, actor, request) { const before = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: true } }); if (!before) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } const integration = await prisma.integration.update({ where: { id: integrationId }, data: { workspaceName: payload.workspaceName ?? before.workspaceName, externalProjectKey: payload.externalProjectKey ?? before.externalProjectKey, baseUrl: payload.baseUrl ?? before.baseUrl, authMode: payload.authMode ?? before.authMode, connectionStatus: payload.connectionStatus ?? before.connectionStatus, syncDirection: payload.syncDirection ?? before.syncDirection, liveEnabled: payload.liveEnabled !== undefined ? payload.liveEnabled : before.liveEnabled, liveConfig: payload.liveConfig ?? before.liveConfig } }); if (payload.credentials) { await saveIntegrationCredentials({ integrationId, credentials: payload.credentials, actor, request }); } return integration; } async function deleteIntegration(integrationId, actor, request) { const before = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: true } }); if (!before) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } await prisma.oAuthToken.deleteMany({ where: { integrationId } }).catch(() => {}); await prisma.integration.delete({ where: { id: integrationId } }); await writeAuditEvent({ organizationId: before.workspace.organizationId, workspaceId: before.workspaceId, actorUserId: actor?.id || null, entityType: "Integration", entityId: integrationId, action: "delete", beforeJson: before, request }); return { deleted: true }; } async function syncIntegration(integrationId, actor, request) { const integration = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: { include: { projects: { include: { risks: { include: { mitigations: true } } } } } } } }); if (!integration) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } const startedAt = new Date(); let createdCount = 0; let updatedCount = 0; let failedCount = 0; const failureLog = []; const statusListMap = integration.liveConfig?.statusListMap || {}; for (const project of integration.workspace.projects) { for (const rawRisk of project.risks) { const risk = scoreRiskForSync(rawRisk); try { const externalKey = `${integration.provider}-${risk.id}`; const existingMapping = await prisma.externalWorkItemMapping.findFirst({ where: { integrationId, localEntityId: risk.id, projectId: project.id } }); const externalStatus = riskStatusToExternalStatus(integration.provider, risk); let fieldMapping = { mode: "simulated" }; if (integration.provider === "TRELLO") { const simulatedListId = statusListMap[risk.status] || integration.externalProjectKey || ""; fieldMapping.simulatedListId = simulatedListId; fieldMapping.simulatedLabels = [ `COSMIC:${risk.dimension}`, `COSMIC:Phase:${risk.lifecyclePhase}` ]; } if (existingMapping) { await prisma.externalWorkItemMapping.update({ where: { id: existingMapping.id }, data: { externalStatus, syncStatus: "Synced", lastSyncedAt: new Date(), fieldMapping } }); updatedCount++; } else { await prisma.externalWorkItemMapping.create({ data: { integrationId, projectId: project.id, localEntityType: "Risk", localEntityId: risk.id, localTitle: risk.title, externalItemType: integration.provider === "TRELLO" ? "Card" : integration.provider === "JIRA" ? "Issue" : "Task", externalItemId: externalKey, externalItemKey: externalKey, externalStatus, syncStatus: "Synced", lastSyncedAt: new Date(), fieldMapping } }); createdCount++; } } catch (err) { failedCount++; failureLog.push({ riskId: risk.id, error: err.message }); } } } const finishedAt = new Date(); const syncRun = await prisma.integrationSyncRun.create({ data: { integrationId, projectId: integration.workspace.projects[0]?.id || "", provider: integration.provider, status: failedCount > 0 ? "Partial" : "Completed", startedAt, finishedAt, createdCount, updatedCount, failedCount, summary: `${integration.provider} sync completed: ${createdCount} created, ${updatedCount} updated, ${failedCount} failed.${integration.provider === "TRELLO" && Object.keys(statusListMap).length ? ` Mapped using ${Object.keys(statusListMap).length} status→list rules.` : ""}`, failureLog } }); await prisma.integration.update({ where: { id: integrationId }, data: { connectionStatus: "CONNECTED", lastSyncAt: new Date() } }); const mappings = await prisma.externalWorkItemMapping.findMany({ where: { integrationId } }); await writeAuditEvent({ organizationId: integration.workspace.organizationId, workspaceId: integration.workspaceId, actorUserId: actor?.id || null, entityType: "IntegrationSyncRun", entityId: syncRun.id, action: "sync", afterJson: { status: syncRun.status, createdCount, updatedCount, failedCount }, request }); return { syncRun, mappings, integration }; } async function testLiveConnector(integrationId, actor, request) { const integration = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: true } }); if (!integration) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } if (!liveModeEnabled(integration)) { const error = new Error("Live PM integration is disabled. Set integration.liveEnabled=true or COSMIC_LIVE_PM_ENABLED=true."); error.status = 409; throw error; } const credentials = await loadIntegrationCredentials(integration); const client = createPmClient(integration.provider, credentials); const result = await client.testConnection(); await prisma.integration.update({ where: { id: integrationId }, data: { connectionStatus: "CONNECTED" } }); await writeAuditEvent({ organizationId: integration.workspace.organizationId, workspaceId: integration.workspaceId, actorUserId: actor?.id || null, entityType: "Integration", entityId: integrationId, action: "live-test", afterJson: { provider: integration.provider, account: result.account }, request }).catch(() => {}); return Object.assign({}, result, { liveEnabled: true, integration: integrationClientView(integration) }); } async function liveSync(integrationId, actor, request) { const integration = await prisma.integration.findUnique({ where: { id: integrationId }, include: { workspace: { include: { projects: { include: { risks: { include: { mitigations: true } } } } } } } }); if (!integration) { const error = new Error(`Integration not found: ${integrationId}`); error.status = 404; throw error; } if (!liveModeEnabled(integration)) { const error = new Error("Live PM integration is disabled. Set integration.liveEnabled=true or COSMIC_LIVE_PM_ENABLED=true."); error.status = 409; throw error; } const credentials = await loadIntegrationCredentials(integration); const client = createPmClient(integration.provider, credentials); const startedAt = new Date(); let createdCount = 0; let updatedCount = 0; let failedCount = 0; const failureLog = []; for (const project of integration.workspace.projects) { for (const rawRisk of project.risks) { const risk = scoreRiskForSync(rawRisk); try { const existingMapping = await prisma.externalWorkItemMapping.findFirst({ where: { integrationId, localEntityId: risk.id, projectId: project.id } }); const result = existingMapping ? await client.updateRiskWorkItem(existingMapping, risk, rawRisk.mitigations || [], integration) : await client.createRiskWorkItem(risk, rawRisk.mitigations || [], integration); const externalStatus = result.externalStatus || riskStatusToExternalStatus(integration.provider, risk); if (existingMapping) { await prisma.externalWorkItemMapping.update({ where: { id: existingMapping.id }, data: { localTitle: risk.title, externalItemId: result.externalId || existingMapping.externalItemId, externalItemKey: result.externalKey || existingMapping.externalItemKey, externalUrl: result.externalUrl || existingMapping.externalUrl, externalStatus, syncStatus: "Live synced", lastSyncedAt: new Date() } }); updatedCount++; } else { await prisma.externalWorkItemMapping.create({ data: { integrationId, projectId: project.id, localEntityType: "Risk", localEntityId: risk.id, localTitle: risk.title, externalItemType: integration.provider === "TRELLO" ? "Card" : integration.provider === "JIRA" ? "Issue" : "Task", externalItemId: result.externalId, externalItemKey: result.externalKey || result.externalId, externalUrl: result.externalUrl || null, externalStatus, syncStatus: "Live synced", lastSyncedAt: new Date(), fieldMapping: { source: "live-api" } } }); createdCount++; } } catch (err) { failedCount++; failureLog.push({ riskId: risk.id, title: risk.title, error: err.message, status: err.status || null }); } } } const finishedAt = new Date(); const syncRun = await prisma.integrationSyncRun.create({ data: { integrationId, projectId: integration.workspace.projects[0]?.id || "", provider: integration.provider, status: failedCount > 0 ? "Partial" : "Completed", startedAt, finishedAt, createdCount, updatedCount, failedCount, summary: `Live ${integration.provider} sync: ${createdCount} created, ${updatedCount} updated, ${failedCount} failed.`, failureLog } }); await prisma.integration.update({ where: { id: integrationId }, data: { connectionStatus: failedCount > 0 ? "NEEDS_CONFIGURATION" : "CONNECTED", lastSyncAt: new Date() } }); const mappings = await prisma.externalWorkItemMapping.findMany({ where: { integrationId } }); await writeAuditEvent({ organizationId: integration.workspace.organizationId, workspaceId: integration.workspaceId, actorUserId: actor?.id || null, entityType: "IntegrationSyncRun", entityId: syncRun.id, action: "liveSync", afterJson: { status: syncRun.status, createdCount, updatedCount, failedCount }, request }); return { syncRun, mappings, integration: integrationClientView(integration) }; } module.exports = { getIntegration, updateIntegration, deleteIntegration, syncIntegration, testLiveConnector, liveSync, saveIntegrationCredentials, loadIntegrationCredentials };