Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
2 changes: 1 addition & 1 deletion packages/core/shared/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@activepieces/shared",
"version": "0.141.0",
"version": "0.142.0",
"type": "commonjs",
"sideEffects": false,
"main": "./dist/src/index.js",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,7 @@ export const Platform = z.object({
ssoDomainVerification: Nullable(SsoDomainVerification),
federatedAuthProviders: FederatedAuthnProviderConfig,
emailAuthEnabled: z.boolean(),
autoCreatePersonalProjects: z.boolean(),
pinnedPieces: z.array(z.string()),
pieceSelectorConfig: Nullable(PieceSelectorConfig),
})
Expand Down Expand Up @@ -233,6 +234,7 @@ export const PlatformWithoutSensitiveData = z.object({
ssoDomain: Nullable(z.string()),
ssoDomainVerification: Nullable(SsoDomainVerification),
emailAuthEnabled: z.boolean(),
autoCreatePersonalProjects: z.boolean(),
pinnedPieces: z.array(z.string()),
pieceSelectorConfig: Nullable(PieceSelectorConfig),
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ export const UpdatePlatformRequestBody = z.object({
cloudAuthEnabled: OptionalBooleanFromQuery,
googleAuthEnabled: OptionalBooleanFromQuery,
emailAuthEnabled: OptionalBooleanFromQuery,
autoCreatePersonalProjects: OptionalBooleanFromQuery,
allowedAuthDomains: OptionalArrayFromQuery(z.string()),
enforceAllowedAuthDomains: OptionalBooleanFromQuery,
pinnedPieces: OptionalArrayFromQuery(z.string()),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,14 +41,6 @@ export const authenticationUtils = (log: FastifyBaseLogger) => ({
const project = isNil(params.projectId)
? findPersonalProject(projects, params.userId) ?? projects?.[0]
: projects.find((project) => project.id === params.projectId)
if (isNil(project)) {
throw new ActivepiecesError({
code: ErrorCode.INVITATION_ONLY_SIGN_UP,
params: {
message: 'No project found for user',
},
})
}
const identity = await userIdentityService(log).getOneOrFail({ id: user.identityId })
if (!identity.verified) {
throw new ActivepiecesError({
Expand Down Expand Up @@ -83,7 +75,7 @@ export const authenticationUtils = (log: FastifyBaseLogger) => ({
newsLetter: identity.newsLetter,
verified: identity.verified,
token,
projectId: project.id,
projectId: project?.id ?? null,
}
},

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import { QueryRunner } from 'typeorm'
import { Migration } from '../../migration'

export class AddAutoCreatePersonalProjectsToPlatform1834000000000 implements Migration {
name = 'AddAutoCreatePersonalProjectsToPlatform1834000000000'
breaking = false
release = '0.87.1'

public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`
ALTER TABLE "platform"
ADD "autoCreatePersonalProjects" boolean NOT NULL DEFAULT true
`)
}

public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`
ALTER TABLE "platform" DROP COLUMN "autoCreatePersonalProjects"
`)
}
}
2 changes: 2 additions & 0 deletions packages/server/api/src/app/database/postgres-connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,7 @@ import { AddAiProviderScopes1830000000000 } from './migration/postgres/183000000
import { AddChatPersonalization1831000000000 } from './migration/postgres/1831000000000-AddChatPersonalization'
import { BackfillChatPersonalizationForExistingUsers1832000000000 } from './migration/postgres/1832000000000-BackfillChatPersonalizationForExistingUsers'
import { ClearRoleFromCompanyPersonalization1833000000000 } from './migration/postgres/1833000000000-ClearRoleFromCompanyPersonalization'
import { AddAutoCreatePersonalProjectsToPlatform1834000000000 } from './migration/postgres/1834000000000-AddAutoCreatePersonalProjectsToPlatform'

const getSslConfig = (): boolean | TlsOptions => {
const useSsl = system.get(AppSystemProp.POSTGRES_USE_SSL)
Expand Down Expand Up @@ -865,6 +866,7 @@ export const getMigrations = (): (new () => Migration)[] => {
AddChatPersonalization1831000000000,
BackfillChatPersonalizationForExistingUsers1832000000000,
ClearRoleFromCompanyPersonalization1833000000000,
AddAutoCreatePersonalProjectsToPlatform1834000000000,
]
return migrations
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -142,14 +142,6 @@ async function assertProjectIsSafeToDelete(projectId: string, callerPlatformId:
},
})
}
if (project.type === ProjectType.PERSONAL) {
throw new ActivepiecesError({
code: ErrorCode.VALIDATION,
params: {
message: 'Personal projects cannot be deleted',
},
})
}
}

async function assertMaximumNumberOfProjectsReachedByEdition(platformId: string, log: FastifyBaseLogger): Promise<void> {
Expand Down
23 changes: 17 additions & 6 deletions packages/server/api/src/app/pieces/piece-sync-service.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { groupBy, tryCatch } from '@activepieces/core-utils'
import { apVersionUtil } from '@activepieces/server-utils'
import { PieceSyncMode, PieceType } from '@activepieces/shared'
import { PieceAudienceFilter, PieceSyncMode, PieceType } from '@activepieces/shared'
import { FastifyBaseLogger } from 'fastify'
import semver from 'semver'
import { rejectedPromiseHandler } from '../helper/promise-handler'
Expand Down Expand Up @@ -52,11 +52,12 @@ export const pieceSyncService = (log: FastifyBaseLogger) => ({
},
}), listCloudPieces()])
log.info({ dbCount: dbPieces.length, cloudCount: cloudPieces.length }, 'Fetched pieces from DB and Cloud')
const added = await installNewPieces(cloudPieces, dbPieces, log, publishCacheRefresh)
const { added, fetchFailed } = await installNewPieces(cloudPieces, dbPieces, log, publishCacheRefresh)
const deleted = await deletePiecesIfNotOnCloud(dbPieces, cloudPieces, log)

log.info({
added,
fetchFailed,
deleted,
durationMs: Math.floor(performance.now() - startTime),
}, 'Piece synchronization completed')
Expand All @@ -83,17 +84,24 @@ async function deletePiecesIfNotOnCloud(dbPieces: PieceMetadataOnly[], cloudPiec
return piecesToDelete.length
}

async function installNewPieces(cloudPieces: PieceRegistryResponse[], dbPieces: PieceMetadataOnly[], log: FastifyBaseLogger, _publishCacheRefresh: boolean): Promise<number> {
async function installNewPieces(cloudPieces: PieceRegistryResponse[], dbPieces: PieceMetadataOnly[], log: FastifyBaseLogger, _publishCacheRefresh: boolean): Promise<{ added: number, fetchFailed: number }> {
const dbMap = new Map<string, true>(dbPieces.map(dbPiece => [`${dbPiece.name}:${dbPiece.version}`, true]))
const newPiecesToFetch = cloudPieces.filter(piece => !dbMap.has(`${piece.name}:${piece.version}`))
const batchSize = 5
let added = 0
let fetchFailed = 0
for (let done = 0; done < newPiecesToFetch.length; done += batchSize) {
const currentBatch = newPiecesToFetch.slice(done, done + batchSize)
await Promise.all(currentBatch.map(async (piece) => {
const url = `${CLOUD_API_URL}/${piece.name}${piece.version ? '?version=' + piece.version : ''}`
const queryParams = new URLSearchParams({ audience: PieceAudienceFilter.ALL })
if (piece.version) {
queryParams.append('version', piece.version)
}
const url = `${CLOUD_API_URL}/${piece.name}?${queryParams.toString()}`
const response = await fetch(url)
if (!response.ok) {
log.warn({ piece: { name: piece.name, version: piece.version }, status: response.status }, '[pieceSyncService#installNewPieces] Error reading piece metadata')
fetchFailed++
return
}
const pieceMetadata = await response.json()
Expand All @@ -106,12 +114,15 @@ async function installNewPieces(cloudPieces: PieceRegistryResponse[], dbPieces:
if (error) {
log.debug({ piece: { name: piece.name, version: piece.version } }, '[pieceSyncService#installNewPieces] Piece already exists, skipping')
}
else {
added++
}
}))
}
if (newPiecesToFetch.length > 0) {
if (added > 0) {
await pieceCache(log).invalidate()
}
return newPiecesToFetch.length
return { added, fetchFailed }
}


Expand Down
5 changes: 5 additions & 0 deletions packages/server/api/src/app/platform/platform.entity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,11 @@ export const PlatformEntity = new EntitySchema<PlatformSchema>({
type: Boolean,
nullable: false,
},
autoCreatePersonalProjects: {
type: Boolean,
nullable: false,
default: true,
},
federatedAuthProviders: {
type: 'jsonb',
select: false,
Expand Down
2 changes: 2 additions & 0 deletions packages/server/api/src/app/platform/platform.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ export const platformService = (log: FastifyBaseLogger) => ({
fullLogoUrl: fullLogoUrl ?? defaultTheme.logos.fullLogoUrl,
favIconUrl: favIconUrl ?? defaultTheme.logos.favIconUrl,
emailAuthEnabled: true,
autoCreatePersonalProjects: true,
enforceAllowedAuthDomains: false,
allowedAuthDomains: [],
federatedAuthProviders: { saml: null },
Expand Down Expand Up @@ -176,6 +177,7 @@ export const platformService = (log: FastifyBaseLogger) => ({
...spreadIfDefined('cloudAuthEnabled', params.cloudAuthEnabled),
...spreadIfDefined('googleAuthEnabled', params.googleAuthEnabled),
...spreadIfDefined('emailAuthEnabled', params.emailAuthEnabled),
...spreadIfDefined('autoCreatePersonalProjects', params.autoCreatePersonalProjects),
...spreadIfDefined(
'enforceAllowedAuthDomains',
params.enforceAllowedAuthDomains,
Expand Down
15 changes: 9 additions & 6 deletions packages/server/api/src/app/user/user-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,12 +45,15 @@ export const userService = (log: FastifyBaseLogger) => ({
platformRole: PlatformRole.MEMBER,
})

await projectService(log).create({
displayName: identity.firstName + '\'s Project',
ownerId: newUser.id,
platformId,
type: ProjectType.PERSONAL,
})
const platform = await platformService(log).getOneOrThrow(platformId)
if (platform.autoCreatePersonalProjects) {
await projectService(log).create({
displayName: identity.firstName + '\'s Project',
ownerId: newUser.id,
platformId,
type: ProjectType.PERSONAL,
})
}
return newUser
}
return user
Expand Down
1 change: 1 addition & 0 deletions packages/server/api/test/helpers/mocks/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,7 @@ export const createMockPlatform = (platform?: Partial<Platform>): Platform => {
logoIconUrl: platform?.logoIconUrl ?? faker.image.urlPlaceholder(),
fullLogoUrl: platform?.fullLogoUrl ?? faker.image.urlPlaceholder(),
emailAuthEnabled: platform?.emailAuthEnabled ?? faker.datatype.boolean(),
autoCreatePersonalProjects: platform?.autoCreatePersonalProjects ?? true,
pinnedPieces: platform?.pinnedPieces ?? [],
favIconUrl: platform?.favIconUrl ?? faker.image.urlPlaceholder(),
cloudAuthEnabled: platform?.cloudAuthEnabled ?? faker.datatype.boolean(),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
import { ActionBase, Audience } from '@activepieces/pieces-framework'
import { PackageType, PieceType } from '@activepieces/shared'
import { FastifyBaseLogger, FastifyInstance } from 'fastify'
import { StatusCodes } from 'http-status-codes'
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
import { pieceCache } from '../../../../src/app/pieces/metadata/piece-cache'
import { pieceRepos } from '../../../../src/app/pieces/metadata/piece-metadata-service'
import { pieceSyncService } from '../../../../src/app/pieces/piece-sync-service'
import { createMockPieceMetadata } from '../../../helpers/mocks'
import { createTestContext } from '../../../helpers/test-context'
import { setupTestEnvironment, teardownTestEnvironment } from '../../../helpers/test-setup'

const originalSyncMode = vi.hoisted(() => {
const previous = process.env.AP_PIECES_SYNC_MODE
process.env.AP_PIECES_SYNC_MODE = 'OFFICIAL_AUTO'
return previous
})

const CLOUD_PIECES_URL = 'https://cloud.activepieces.com/api/v1/pieces'
const PIECE_NAME = '@activepieces/piece-audience-sync-probe'
const PIECE_VERSION = '1.0.0'

let app: FastifyInstance
let mockLog: FastifyBaseLogger

function mockAction(name: string, audience: Audience): ActionBase {
return {
name,
displayName: name,
description: `${name} description`,
props: {},
requireAuth: false,
audience,
}
}

const cloudPiece = createMockPieceMetadata({
name: PIECE_NAME,
version: PIECE_VERSION,
pieceType: PieceType.OFFICIAL,
packageType: PackageType.REGISTRY,
actions: {
human_action: mockAction('human_action', 'human'),
ai_action: mockAction('ai_action', 'ai'),
shared_action: mockAction('shared_action', 'both'),
},
})

function applyCloudAudienceFilter(audience: string | null): typeof cloudPiece {
if (audience === 'all') {
return cloudPiece
}
return {
...cloudPiece,
actions: Object.fromEntries(
Object.entries(cloudPiece.actions).filter(([, action]) => action.audience !== 'ai'),
),
}
}

function jsonResponse(body: unknown): Response {
return new Response(JSON.stringify(body), {
status: 200,
headers: { 'content-type': 'application/json' },
})
}

function stubCloudFetch(): void {
const realFetch = globalThis.fetch
vi.stubGlobal('fetch', async (input: string | URL | Request, init?: RequestInit) => {
const target = input instanceof Request ? input.url : String(input)
if (!target.startsWith(CLOUD_PIECES_URL)) {
return realFetch(input, init)
}
const url = new URL(target)
if (url.pathname === '/api/v1/pieces/registry') {
return jsonResponse([{ name: PIECE_NAME, version: PIECE_VERSION }])
}
if (url.pathname === `/api/v1/pieces/${PIECE_NAME}`) {
return jsonResponse(applyCloudAudienceFilter(url.searchParams.get('audience')))
}
return jsonResponse({ message: 'not found' })
})
}

// setup() fires an unawaited boot sync; let it finish before tests truncate and
// re-sync, so it can't race them or outlive the fetch stub.
async function settleBootSync(): Promise<void> {
const deadline = Date.now() + 15_000
while (Date.now() < deadline) {
const row = await pieceRepos().findOneBy({ name: PIECE_NAME, version: PIECE_VERSION })
if (row !== null) {
return
}
await new Promise(resolve => setTimeout(resolve, 100))
}
throw new Error('Boot piece sync did not settle in time')
}

beforeAll(async () => {
stubCloudFetch()
app = await setupTestEnvironment({ fresh: true })
mockLog = app.log
await settleBootSync()
})

afterAll(async () => {
vi.unstubAllGlobals()
if (originalSyncMode === undefined) {
delete process.env.AP_PIECES_SYNC_MODE
}
else {
process.env.AP_PIECES_SYNC_MODE = originalSyncMode
}
await teardownTestEnvironment()
})

beforeEach(async () => {
await pieceRepos().createQueryBuilder().delete().execute()
})

describe('Piece Sync Audience', () => {
it('stores every action of a synced piece, including audience ai', async () => {
await pieceSyncService(mockLog).sync({ publishCacheRefresh: false })

const stored = await pieceRepos().findOneByOrFail({
name: PIECE_NAME,
version: PIECE_VERSION,
})

expect(Object.keys(stored.actions).sort()).toEqual(['ai_action', 'human_action', 'shared_action'])
})

it('keeps ai actions hidden from the default read paths', async () => {
await pieceSyncService(mockLog).sync({ publishCacheRefresh: false })
await pieceCache(mockLog).setup()
const ctx = await createTestContext(app)

const getResponse = await ctx.get(`/v1/pieces/${PIECE_NAME}`)
expect(getResponse.statusCode).toBe(StatusCodes.OK)
const piece: { actions: Record<string, ActionBase> } = getResponse.json()
expect(Object.keys(piece.actions).sort()).toEqual(['human_action', 'shared_action'])

const listResponse = await ctx.get('/v1/pieces')
expect(listResponse.statusCode).toBe(StatusCodes.OK)
const summaries: { name: string, actions: number }[] = listResponse.json()
const summary = summaries.find(item => item.name === PIECE_NAME)
expect(summary?.actions).toBe(2)
})
})
Loading
Loading