Skip to content

Commit aa220d7

Browse files
authored
Merge pull request #360 from RhysSullivan/rs/connections-mcp-migrate
connections: migrate legacy MCP + google-discovery oauth rows
2 parents 22a8c1b + e78249c commit aa220d7

5 files changed

Lines changed: 777 additions & 138 deletions

File tree

Lines changed: 332 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,332 @@
1+
// ---------------------------------------------------------------------------
2+
// MCP OAuth legacy → Connection backfill (cloud)
3+
// ---------------------------------------------------------------------------
4+
//
5+
// Companion to migrate-connections.ts (which handled OpenAPI). Same shape,
6+
// different table: scans `mcp_source`, finds rows whose `config.auth` is
7+
// still on the pre-refactor inline-OAuth shape, mints a stable per-source
8+
// Connection row, backfills the secret routing rows, and rewrites
9+
// `config.auth` to the `{kind:"oauth2", connectionId}` pointer.
10+
//
11+
// Dry-run by default. `--apply` runs the per-row transactions.
12+
//
13+
// Run (dry-run):
14+
// op run --env-file=.env.production -- bun run scripts/migrate-mcp-connections.ts
15+
// Run (apply):
16+
// op run --env-file=.env.production -- bun run scripts/migrate-mcp-connections.ts --apply
17+
18+
import { Option, Schema } from "effect";
19+
import postgres from "postgres";
20+
21+
import { McpConnectionAuth } from "@executor/plugin-mcp";
22+
23+
const APPLY = process.argv.includes("--apply");
24+
const DUMP_BLOCKED = process.argv.includes("--dump-blocked");
25+
26+
// ---------------------------------------------------------------------------
27+
// Legacy MCP oauth2 auth shape (pre-refactor). Inlined — this script is the
28+
// only place that still needs to know about it. Once cloud + local have
29+
// run this migration, the file can be deleted.
30+
// ---------------------------------------------------------------------------
31+
32+
const JsonObject = Schema.Record({ key: Schema.String, value: Schema.Unknown });
33+
34+
const LegacyMcpOAuth2 = Schema.Struct({
35+
kind: Schema.Literal("oauth2"),
36+
accessTokenSecretId: Schema.String,
37+
refreshTokenSecretId: Schema.NullOr(Schema.String),
38+
tokenType: Schema.optionalWith(Schema.String, { default: () => "Bearer" }),
39+
expiresAt: Schema.NullOr(Schema.Number),
40+
scope: Schema.NullOr(Schema.String),
41+
clientInformation: Schema.optionalWith(Schema.NullOr(JsonObject), {
42+
default: () => null,
43+
}),
44+
authorizationServerUrl: Schema.optionalWith(Schema.NullOr(Schema.String), {
45+
default: () => null,
46+
}),
47+
resourceMetadataUrl: Schema.optionalWith(Schema.NullOr(Schema.String), {
48+
default: () => null,
49+
}),
50+
});
51+
type LegacyMcpOAuth2 = typeof LegacyMcpOAuth2.Type;
52+
53+
const decodeCurrentAuth = Schema.decodeUnknownOption(McpConnectionAuth);
54+
const decodeLegacy = Schema.decodeUnknownOption(LegacyMcpOAuth2);
55+
56+
const isRecord = (v: unknown): v is Record<string, unknown> =>
57+
typeof v === "object" && v !== null && !Array.isArray(v);
58+
59+
// ---------------------------------------------------------------------------
60+
// Row classification
61+
// ---------------------------------------------------------------------------
62+
63+
type Row = {
64+
scope_id: string;
65+
id: string;
66+
name: string;
67+
config: unknown;
68+
};
69+
70+
type Bucket =
71+
| { kind: "non-remote"; row: Row }
72+
| { kind: "no-oauth"; row: Row }
73+
| { kind: "current"; row: Row }
74+
| { kind: "legacy-migratable"; row: Row; legacy: LegacyMcpOAuth2; endpoint: string }
75+
| { kind: "legacy-blocked"; row: Row; legacy: LegacyMcpOAuth2; reason: string }
76+
| { kind: "unknown"; row: Row; shape: string };
77+
78+
const classifyRow = (row: Row): Bucket => {
79+
if (!isRecord(row.config)) return { kind: "unknown", row, shape: typeof row.config };
80+
if (row.config.transport !== "remote") return { kind: "non-remote", row };
81+
const auth = row.config.auth;
82+
if (!isRecord(auth)) return { kind: "no-oauth", row };
83+
if (auth.kind !== "oauth2") return { kind: "no-oauth", row };
84+
85+
if (Option.isSome(decodeCurrentAuth(auth))) return { kind: "current", row };
86+
87+
const legacyOption = decodeLegacy(auth);
88+
if (Option.isSome(legacyOption)) {
89+
const legacy = legacyOption.value;
90+
const endpoint =
91+
typeof row.config.endpoint === "string" ? row.config.endpoint : null;
92+
if (!endpoint) {
93+
return {
94+
kind: "legacy-blocked",
95+
row,
96+
legacy,
97+
reason: "config.endpoint missing",
98+
};
99+
}
100+
return { kind: "legacy-migratable", row, legacy, endpoint };
101+
}
102+
103+
const shape = `{${Object.keys(auth).sort().join(",")}}`;
104+
return { kind: "unknown", row, shape };
105+
};
106+
107+
// ---------------------------------------------------------------------------
108+
// Main
109+
// ---------------------------------------------------------------------------
110+
111+
type SecretRow = {
112+
id: string;
113+
scope_id: string;
114+
owned_by_connection_id: string | null;
115+
};
116+
117+
const main = async () => {
118+
const connectionString =
119+
process.env.DATABASE_URL || process.env.HYPERDRIVE_CONNECTION_STRING || "";
120+
if (!connectionString) {
121+
console.error(
122+
"DATABASE_URL not set (try: op run --env-file=.env.production -- ...)",
123+
);
124+
process.exit(1);
125+
}
126+
127+
const sql = postgres(connectionString, {
128+
max: 1,
129+
onnotice: () => undefined,
130+
ssl: "require",
131+
});
132+
133+
try {
134+
const rows = (await sql<Row[]>`
135+
select scope_id, id, name, config
136+
from mcp_source
137+
`) as Row[];
138+
139+
console.log(`\nScanned ${rows.length} mcp_source row(s)`);
140+
console.log(APPLY ? "Mode: APPLY (writes enabled)\n" : "Mode: DRY-RUN (no writes)\n");
141+
142+
const buckets = rows.map(classifyRow);
143+
144+
const counts = {
145+
"non-remote": 0,
146+
"no-oauth": 0,
147+
current: 0,
148+
"legacy-migratable": 0,
149+
"legacy-blocked": 0,
150+
unknown: 0,
151+
};
152+
for (const b of buckets) counts[b.kind]++;
153+
154+
console.log("Classification:");
155+
console.log(` non-remote (stdio): ${counts["non-remote"]}`);
156+
console.log(` no oauth2 auth: ${counts["no-oauth"]}`);
157+
console.log(` already on new shape: ${counts.current}`);
158+
console.log(` legacy — would migrate: ${counts["legacy-migratable"]}`);
159+
console.log(` legacy — blocked: ${counts["legacy-blocked"]}`);
160+
console.log(` unrecognized shape: ${counts.unknown}\n`);
161+
162+
const migratable = buckets.filter(
163+
(b): b is Extract<Bucket, { kind: "legacy-migratable" }> =>
164+
b.kind === "legacy-migratable",
165+
);
166+
const blocked = buckets.filter(
167+
(b) => b.kind === "legacy-blocked" || b.kind === "unknown",
168+
);
169+
170+
for (const b of blocked) {
171+
const ref = `${b.row.scope_id}/${b.row.id}`;
172+
if (b.kind === "legacy-blocked") {
173+
console.log(`[BLOCKED] ${ref}`);
174+
console.log(` reason: ${b.reason}`);
175+
} else if (b.kind === "unknown") {
176+
console.log(`[UNKNOWN] ${ref}`);
177+
console.log(` shape: ${b.shape}`);
178+
}
179+
if (DUMP_BLOCKED && isRecord(b.row.config)) {
180+
console.log(` auth: ${JSON.stringify(b.row.config.auth, null, 2)}`);
181+
}
182+
console.log();
183+
}
184+
185+
for (const b of migratable) {
186+
const ref = `${b.row.scope_id}/${b.row.id}`;
187+
console.log(`[migratable] ${ref}`);
188+
console.log(` endpoint: ${b.endpoint}`);
189+
console.log(` accessTokenSecret: ${b.legacy.accessTokenSecretId}`);
190+
console.log(
191+
` refreshTokenSecret: ${b.legacy.refreshTokenSecretId ?? "(null)"}`,
192+
);
193+
console.log(
194+
` authServerUrl: ${b.legacy.authorizationServerUrl ?? "(null)"}`,
195+
);
196+
console.log();
197+
}
198+
199+
if (blocked.length > 0) {
200+
console.log(
201+
`ABORT: ${blocked.length} row(s) blocked or unrecognized; inspect above and fix or delete before re-running.`,
202+
);
203+
process.exit(2);
204+
}
205+
206+
if (!APPLY) {
207+
console.log(`OK: ${migratable.length} row(s) would migrate. Re-run with --apply.`);
208+
return;
209+
}
210+
211+
if (migratable.length === 0) {
212+
console.log("Nothing to migrate.");
213+
return;
214+
}
215+
216+
let applied = 0;
217+
let failed = 0;
218+
for (const b of migratable) {
219+
const ref = `${b.row.scope_id}/${b.row.id}`;
220+
const connectionId = `mcp-oauth2-${b.row.id}`;
221+
const l = b.legacy;
222+
const providerState = {
223+
endpoint: b.endpoint,
224+
tokenType: l.tokenType,
225+
clientInformation: l.clientInformation,
226+
authorizationServerUrl: l.authorizationServerUrl,
227+
authorizationServerMetadata: null,
228+
resourceMetadataUrl: l.resourceMetadataUrl,
229+
resourceMetadata: null,
230+
};
231+
const authPointer = { kind: "oauth2" as const, connectionId };
232+
const nextConfig = { ...(isRecord(b.row.config) ? b.row.config : {}), auth: authPointer };
233+
234+
try {
235+
await sql.begin(async (tx) => {
236+
await tx`
237+
insert into connection (
238+
id, scope_id, provider, kind, identity_label,
239+
access_token_secret_id, refresh_token_secret_id,
240+
expires_at, scope, provider_state,
241+
created_at, updated_at
242+
) values (
243+
${connectionId},
244+
${b.row.scope_id},
245+
${"mcp:oauth2"},
246+
${"user"},
247+
${b.row.name},
248+
${l.accessTokenSecretId},
249+
${l.refreshTokenSecretId},
250+
${l.expiresAt},
251+
${l.scope},
252+
${tx.json(providerState)},
253+
now(),
254+
now()
255+
)
256+
`;
257+
258+
const secretIds = [l.accessTokenSecretId];
259+
if (l.refreshTokenSecretId) secretIds.push(l.refreshTokenSecretId);
260+
261+
const existing = (await tx<SecretRow[]>`
262+
select id, scope_id, owned_by_connection_id
263+
from secret
264+
where scope_id = ${b.row.scope_id} and id = any(${secretIds})
265+
`) as SecretRow[];
266+
const alreadyOwned = existing.filter(
267+
(r) =>
268+
r.owned_by_connection_id !== null &&
269+
r.owned_by_connection_id !== connectionId,
270+
);
271+
if (alreadyOwned.length > 0) {
272+
throw new Error(
273+
`secret(s) already owned: ${alreadyOwned.map((r) => `${r.id}(owner=${r.owned_by_connection_id})`).join(", ")}`,
274+
);
275+
}
276+
// Backfill any missing routing rows pointing at workos-vault, the
277+
// only writable provider on cloud. Matches the openapi migration.
278+
const missing = secretIds.filter(
279+
(id) => !existing.some((r) => r.id === id),
280+
);
281+
for (const id of missing) {
282+
const name =
283+
id === l.accessTokenSecretId
284+
? `Connection ${connectionId} access token`
285+
: `Connection ${connectionId} refresh token`;
286+
await tx`
287+
insert into secret (
288+
id, scope_id, provider, name,
289+
owned_by_connection_id, created_at
290+
) values (
291+
${id}, ${b.row.scope_id}, ${"workos-vault"}, ${name},
292+
${connectionId}, now()
293+
)
294+
`;
295+
}
296+
if (existing.length > 0) {
297+
await tx`
298+
update secret
299+
set owned_by_connection_id = ${connectionId}
300+
where scope_id = ${b.row.scope_id} and id = any(${secretIds})
301+
`;
302+
}
303+
304+
await tx`
305+
update mcp_source
306+
set config = ${tx.json(nextConfig)}
307+
where scope_id = ${b.row.scope_id} and id = ${b.row.id}
308+
`;
309+
});
310+
applied++;
311+
console.log(` [OK] ${ref} -> ${connectionId}`);
312+
} catch (err) {
313+
failed++;
314+
console.log(
315+
` [FAIL] ${ref}: ${err instanceof Error ? err.message : String(err)}`,
316+
);
317+
}
318+
}
319+
320+
console.log();
321+
console.log(`Applied: ${applied}`);
322+
console.log(`Failed: ${failed}`);
323+
if (failed > 0) process.exit(3);
324+
} finally {
325+
await sql.end({ timeout: 5 });
326+
}
327+
};
328+
329+
main().catch((err) => {
330+
console.error("migrate-mcp-connections failed:", err);
331+
process.exit(1);
332+
});

‎apps/local/src/server/executor.ts‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import {
1414
readLegacySecrets,
1515
type LegacySecret,
1616
} from "./db-upgrade";
17-
import { migrateOpenApiOAuthConnections } from "./migrate-connections";
17+
import { migrateLegacyConnections } from "./migrate-connections";
1818

1919
import {
2020
Scope,
@@ -137,9 +137,10 @@ const createLocalExecutorLayer = () => {
137137
if (legacySecrets.length > 0) {
138138
importLegacySecrets(sqlite, scopeId, legacySecrets);
139139
}
140-
// Upgrade pre-connection OpenAPI OAuth rows onto the new Connection
141-
// pointer shape. Idempotent + no-op when there's nothing to migrate.
142-
yield* Effect.promise(() => migrateOpenApiOAuthConnections(sqlite));
140+
// Upgrade pre-connection openapi/mcp/google-discovery rows onto the
141+
// new Connection pointer shape. Idempotent + no-op when there's
142+
// nothing to migrate.
143+
yield* Effect.promise(() => migrateLegacyConnections(sqlite));
143144
const configPath = resolveConfigPath(cwd);
144145
const configFile = makeFileConfigSink({
145146
path: configPath,

0 commit comments

Comments
 (0)