Skip to content
Open
Show file tree
Hide file tree
Changes from 11 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
1 change: 1 addition & 0 deletions .prettierignore
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ pnpm-lock.yaml
# Output
gen
dist
connect/src/wg

# Next.js build output
.next
Expand Down
2,427 changes: 1,271 additions & 1,156 deletions connect-go/gen/proto/wg/cosmo/platform/v1/platform.pb.go

Large diffs are not rendered by default.

297 changes: 177 additions & 120 deletions connect/src/wg/cosmo/platform/v1/platform_pb.ts

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
import { HandlerContext } from '@connectrpc/connect';
import { EnumStatusCode } from '@wundergraph/cosmo-connect/dist/common/common_pb';
import {
FeatureSubgraphInFlagComposition,
GetFeatureFlagsInLatestCompositionByFederatedGraphRequest,
GetFeatureFlagsInLatestCompositionByFederatedGraphResponse,
} from '@wundergraph/cosmo-connect/dist/platform/v1/platform_pb';
import { PlainMessage, FeatureFlagDTO } from '../../../types/index.js';
import { FeatureFlagRepository } from '../../repositories/FeatureFlagRepository.js';
import { FederatedGraphRepository } from '../../repositories/FederatedGraphRepository.js';
import { GraphCompositionRepository } from '../../repositories/GraphCompositionRepository.js';
import { NamespaceRepository } from '../../repositories/NamespaceRepository.js';
import type { RouterOptions } from '../../routes.js';
import { enrichLogger, getLogger, handleError } from '../../util.js';
Expand Down Expand Up @@ -37,6 +39,7 @@ export function getFeatureFlagsInLatestCompositionByFederatedGraph(
details: `Namespace ${req.namespace} not found`,
},
featureFlags: [],
featureSubgraphs: [],
};
}

Expand All @@ -48,6 +51,7 @@ export function getFeatureFlagsInLatestCompositionByFederatedGraph(
details: `Federated Graph '${req.federatedGraphName}' not found`,
},
featureFlags: [],
featureSubgraphs: [],
};
}

Expand All @@ -62,28 +66,54 @@ export function getFeatureFlagsInLatestCompositionByFederatedGraph(
});

const featureFlags: FeatureFlagDTO[] = [];
if (ffsInLatestValidComposition) {
for (const ff of ffsInLatestValidComposition) {
if (!ff.featureFlagId) {
continue;
}
const flag = await featureFlagRepo.getFeatureFlagById({
featureFlagId: ff.featureFlagId,
namespaceId: namespace.id,
includeSubgraphs: false,
});
if (flag) {
// True means the composition reported for this flag is its last successful one, not its latest.
featureFlags.push({ ...flag, hasFailedLatestComposition: ff.hasFailedLatestComposition });
}
const flagIdByComposedSchemaVersionId = new Map<string, string>();
for (const ff of ffsInLatestValidComposition ?? []) {
if (!ff.featureFlagId) {
continue;
}
const flag = await featureFlagRepo.getFeatureFlagById({
featureFlagId: ff.featureFlagId,
namespaceId: namespace.id,
includeSubgraphs: false,
});
if (flag) {
// True means the composition reported for this flag is its last successful one, not its latest.
featureFlags.push({ ...flag, hasFailedLatestComposition: ff.hasFailedLatestComposition });
flagIdByComposedSchemaVersionId.set(ff.id, ff.featureFlagId);
}
}

const compositionRepo = new GraphCompositionRepository(logger, opts.db);
const pinnedFeatureSubgraphs = await compositionRepo.getFeatureSubgraphsByComposedSchemaVersionIds({
schemaVersionIds: [...flagIdByComposedSchemaVersionId.keys()],
organizationId: authContext.organizationId,
rbac: authContext.rbac,
});

const featureSubgraphs: PlainMessage<FeatureSubgraphInFlagComposition>[] = [];
for (const pinned of pinnedFeatureSubgraphs) {
const featureFlagId = flagIdByComposedSchemaVersionId.get(pinned.composedSchemaVersionId);
if (!featureFlagId) {
continue;
}

featureSubgraphs.push({
featureFlagId,
id: pinned.id,
name: pinned.name,
targetId: pinned.targetId,
schemaVersionId: pinned.schemaVersionId,
routingUrl: pinned.routingUrl,
subscriptionUrl: pinned.subscriptionUrl ?? '',
});
}

return {
response: {
code: EnumStatusCode.OK,
},
featureFlags,
featureSubgraphs,
};
},
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,16 @@ import {
graphCompositionSubgraphs,
schemaVersion,
subgraphs,
targets,
users,
} from '../../db/schema.js';
import { DateRange, GraphCompositionDTO } from '../../types/index.js';
import { CompositionSubgraphRecord } from '../composition/composer.js';
import { RBACEvaluator } from '../services/RBACEvaluator.js';
import { traced } from '../tracing.js';
import { FederatedGraphRepository } from './FederatedGraphRepository.js';
import { OrganizationRepository } from './OrganizationRepository.js';
import { SubgraphRepository } from './SubgraphRepository.js';

@traced
export class GraphCompositionRepository {
Expand Down Expand Up @@ -411,6 +414,52 @@ export class GraphCompositionRepository {
return [...compositionSubgraphs, ...childCompositionSubgraphs];
}

/**
* @param input.schemaVersionIds Composed schema versions of the compositions to read.
* @returns A row per feature subgraph per composition, where `schemaVersionId` is the version
* that composition froze rather than the latest published one. Feature subgraphs deleted since
* the composition are omitted, even though `graph_composition_subgraphs` retains their rows.
*/
public async getFeatureSubgraphsByComposedSchemaVersionIds(input: {
schemaVersionIds: string[];
organizationId: string;
rbac?: RBACEvaluator;
}) {
Comment thread
gausie marked this conversation as resolved.
if (input.schemaVersionIds.length === 0) {
return [];
}

const conditions: (SQL<unknown> | undefined)[] = [
inArray(graphCompositions.schemaVersionId, input.schemaVersionIds),
eq(schemaVersion.organizationId, input.organizationId),
eq(graphCompositionSubgraphs.isFeatureSubgraph, true),
not(eq(graphCompositionSubgraphs.changeType, 'removed')),
];
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// The query must join targets, which the RBAC conditions gate on.
if (!SubgraphRepository.applyRbacConditionsToQuery(input.rbac, conditions)) {
return [];
}

return await this.db
.select({
composedSchemaVersionId: graphCompositions.schemaVersionId,
id: graphCompositionSubgraphs.subgraphId,
name: graphCompositionSubgraphs.subgraphName,
targetId: graphCompositionSubgraphs.subgraphTargetId,
schemaVersionId: graphCompositionSubgraphs.schemaVersionId,
routingUrl: subgraphs.routingUrl,
subscriptionUrl: subgraphs.subscriptionUrl,
})
.from(graphCompositionSubgraphs)
.innerJoin(graphCompositions, eq(graphCompositions.id, graphCompositionSubgraphs.graphCompositionId))
.innerJoin(schemaVersion, eq(schemaVersion.id, graphCompositions.schemaVersionId))
.innerJoin(subgraphs, eq(subgraphs.id, graphCompositionSubgraphs.subgraphId))
.innerJoin(targets, eq(targets.id, subgraphs.targetId))
.where(and(...conditions))
.execute();
}

public async getGraphCompositions({
fedGraphTargetId,
organizationId,
Expand Down
Loading
Loading