/home/techb158/balavpn.abdallabala.com/src/services
Edit: /home/techb158/balavpn.abdallabala.com/src/services/integration-service.js (16595B)
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
};