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
16 changes: 12 additions & 4 deletions src/api/auth/identity-mapping-plan.ts
Original file line number Diff line number Diff line change
@@ -1,23 +1,28 @@
import { createHash } from 'node:crypto';
import { identityEndpointSchema } from '@treeseed/sdk/identity';

export interface IdentityMapping { userId: string; issuer: string; subject: string }
export interface IdentityMappingInventory {
users: Array<{ id: string; status: string }>;
mappings: IdentityMapping[];
workloads: Array<{ id: string; issuer: string; subject: string }>;
}

/** Read-only migration planning. Inputs must come from an authenticated import,
* never from email matching. Applying requires a coordinated database restore point. */
export function planIdentityMappings(inventory: IdentityMappingInventory, requested: IdentityMapping[]) {
const normalize = (mapping: IdentityMapping): IdentityMapping => {
const url = new URL(mapping.issuer);
if (url.protocol !== 'https:' || url.username || url.password || url.search || url.hash ||
!mapping.subject.trim() || !mapping.userId.trim()) throw new Error('Invalid explicit identity mapping');
identityEndpointSchema.parse(mapping.issuer);
if (typeof mapping.userId !== 'string' || !/^[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$/u.test(mapping.userId)
|| typeof mapping.subject !== 'string' || !/^[\x21-\x7e]{1,255}$/u.test(mapping.subject))
throw new Error('Invalid explicit identity mapping');
// OIDC issuer identifiers are exact strings; do not normalize trailing slashes.
return { userId: mapping.userId, issuer: mapping.issuer, subject: mapping.subject };
};
const key = (mapping: IdentityMapping) => JSON.stringify([mapping.issuer, mapping.subject]);
const desired = requested.map(normalize).sort((a, b) => key(a).localeCompare(key(b), 'en'));
const workloadIds = new Set(inventory.workloads.map(value => value.id));
const workloadSubjects = new Set(inventory.workloads.map(value => JSON.stringify([value.issuer, value.subject])));
const existing = new Map<string, string>();
for (const mapping of inventory.mappings) {
const previous = existing.get(key(mapping));
Expand All @@ -28,6 +33,7 @@ export function planIdentityMappings(inventory: IdentityMappingInventory, reques
const operations = desired.map(mapping => {
if (seen.has(key(mapping))) throw new Error('Duplicate requested identity');
seen.add(key(mapping));
if (workloadIds.has(mapping.userId) || workloadSubjects.has(key(mapping))) throw new Error('Human mapping conflicts with workload identity');
const users = inventory.users.filter(user => user.id === mapping.userId);
if (users.length !== 1 || users[0].status !== 'active') throw new Error('Mapping requires one active existing user');
const owner = existing.get(key(mapping));
Expand All @@ -38,6 +44,8 @@ export function planIdentityMappings(inventory: IdentityMappingInventory, reques
const observed = {
users: [...inventory.users].sort((a, b) => a.id.localeCompare(b.id, 'en')),
mappings: [...inventory.mappings].sort((a, b) => key(a).localeCompare(key(b), 'en')),
workloads: [...inventory.workloads].sort((a, b) => a.id.localeCompare(b.id, 'en')),
};
return { inventoryDigest: createHash('sha256').update(JSON.stringify(observed)).digest('hex'), operations };
return { inventoryDigest: createHash('sha256').update(JSON.stringify(observed)).digest('hex'),
requestDigest: createHash('sha256').update(JSON.stringify(desired)).digest('hex'), operations };
}
12 changes: 7 additions & 5 deletions src/api/auth/identity-mapping-transaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,17 +7,19 @@ interface MappingDatabase { transaction<T>(run: (client: PoolClient) => Promise<
/** Internal migration primitive, not an account-linking endpoint. The orchestrator
* must authorize the import and establish a coordinated restore point first. */
export async function applyIdentityMappings(database: MappingDatabase, input: {
requested: IdentityMapping[]; inventoryDigest: string;
requested: IdentityMapping[]; inventoryDigest: string; requestDigest: string;
}) {
return database.transaction(async client => {
// Lock both inventory tables, including against legacy inserts/updates. Keep
// Lock human and workload inventory, including legacy inserts/updates. Keep
// this maintenance transaction short; no network calls while locks are held.
await client.query("SET LOCAL lock_timeout = '5s'");
await client.query('LOCK TABLE users, user_identities IN SHARE ROW EXCLUSIVE MODE');
await client.query('LOCK TABLE users, user_identities, identity_workloads IN SHARE ROW EXCLUSIVE MODE');
const users = await client.query('SELECT id, status FROM users');
const mappings = await client.query('SELECT user_id AS "userId", provider AS issuer, provider_subject AS subject FROM user_identities');
const plan = planIdentityMappings({ users: users.rows, mappings: mappings.rows }, input.requested);
if (plan.inventoryDigest !== input.inventoryDigest) throw new Error('Identity inventory changed; create a new mapping plan');
const workloads = await client.query('SELECT id,issuer,subject FROM identity_workloads');
const plan = planIdentityMappings({ users: users.rows, mappings: mappings.rows, workloads: workloads.rows }, input.requested);
if (plan.inventoryDigest !== input.inventoryDigest || plan.requestDigest !== input.requestDigest)
throw new Error('Identity inventory or request changed; create a new mapping plan');
const now = new Date().toISOString();
for (const mapping of plan.operations) {
if (mapping.action === 'noop') continue;
Expand Down
23 changes: 19 additions & 4 deletions tests/unit/control-plane/accounts/identity-mapping-plan.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { describe, expect, it } from 'vitest';
import { planIdentityMappings } from '../../../../src/api/auth/identity-mapping-plan.ts';

const mapping = { userId: 'existing-user', issuer: 'https://identity.example.test/realms/local', subject: 'imported-subject' };
const inventory = { users: [{ id: mapping.userId, status: 'active' }], mappings: [] };
const inventory = { users: [{ id: mapping.userId, status: 'active' }], mappings: [], workloads: [] };
describe('explicit identity migration planning', () => {
it('preserves user IDs and does not mutate the inventory', () => {
expect(planIdentityMappings(inventory, [mapping]).operations).toEqual([{ ...mapping, action: 'bind' }]);
Expand All @@ -16,8 +16,8 @@ describe('explicit identity migration planning', () => {
expect(() => planIdentityMappings(inventory, [mapping, mapping])).toThrow('Duplicate');
});
it('rejects absent and disabled users instead of creating accounts', () => {
expect(() => planIdentityMappings({ users: [], mappings: [] }, [mapping])).toThrow('active existing user');
expect(() => planIdentityMappings({ users: [{ id: mapping.userId, status: 'disabled' }], mappings: [] }, [mapping])).toThrow('active existing user');
expect(() => planIdentityMappings({ ...inventory, users: [] }, [mapping])).toThrow('active existing user');
expect(() => planIdentityMappings({ ...inventory, users: [{ id: mapping.userId, status: 'disabled' }] }, [mapping])).toThrow('active existing user');
});
it('rejects unsafe issuers and empty subjects', () => {
for (const issuer of ['http://identity.test', 'https://user:password@identity.test', 'https://identity.test/?q=1', 'https://identity.test/#fragment']) {
Expand All @@ -27,8 +27,23 @@ describe('explicit identity migration planning', () => {
});
it('binds the plan to account status and preserves exact issuer distinctions', () => {
const before = planIdentityMappings(inventory, []);
const after = planIdentityMappings({ users: [{ id: mapping.userId, status: 'disabled' }], mappings: [] }, []);
const after = planIdentityMappings({ ...inventory, users: [{ id: mapping.userId, status: 'disabled' }] }, []);
expect(before.inventoryDigest).not.toBe(after.inventoryDigest);
expect(planIdentityMappings({ ...inventory, mappings: [mapping] }, [{ ...mapping, issuer: `${mapping.issuer}/` }]).operations[0].action).toBe('bind');
});
it('rejects workload subjects and principal IDs without treating them as human accounts', () => {
for (const workload of [{ id: 'service', issuer: mapping.issuer, subject: mapping.subject },
{ id: mapping.userId, issuer: mapping.issuer, subject: 'service-subject' }]) {
expect(() => planIdentityMappings({ ...inventory, workloads: [workload] }, [mapping])).toThrow('workload identity');
}
});
it('binds exact requests and all workload inventory while retaining valid opaque subjects', () => {
const before = planIdentityMappings(inventory, [mapping]);
const changed = planIdentityMappings(inventory, [{ ...mapping, subject: 'other|subject' }]);
expect(changed.requestDigest).not.toBe(before.requestDigest);
expect(changed.inventoryDigest).toBe(before.inventoryDigest);
expect(planIdentityMappings({ ...inventory, workloads: [{ id: 'service', issuer: mapping.issuer, subject: 'other' }] }, [mapping]).inventoryDigest).not.toBe(before.inventoryDigest);
for (const subject of ['x'.repeat(256), 'bad\nsubject', 'with space'])
expect(() => planIdentityMappings(inventory, [{ ...mapping, subject }])).toThrow();
});
});
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { randomUUID } from 'node:crypto';
import { readFileSync } from 'node:fs';
import pg from 'pg';
import { describe, expect, it } from 'vitest';
import { AUTH_SCHEMA_SQL } from '../../../../src/api/auth/postgres-store.ts';
Expand All @@ -19,6 +20,7 @@ describe.skipIf(!url)('transactional identity mapping in real PostgreSQL', () =>
connection.pathname = `/${name}`;
pool = new pg.Pool({ connectionString: connection.href });
for (const sql of AUTH_SCHEMA_SQL.slice(0, 3)) await pool.query(sql);
await pool.query(readFileSync('drizzle/control-plane/0020_identity_workloads.sql', 'utf8'));
await pool.query("INSERT INTO users(id,status,created_at,updated_at) VALUES ('preserved','active','now','now')");
const database = { transaction: async <T>(run: (client: pg.PoolClient) => Promise<T>) => {
const client = await pool!.connect();
Expand All @@ -27,19 +29,24 @@ describe.skipIf(!url)('transactional identity mapping in real PostgreSQL', () =>
finally { client.release(); }
} };
const inventory = async () => ({ users: (await pool!.query('SELECT id,status FROM users')).rows,
mappings: (await pool!.query('SELECT user_id AS "userId",provider AS issuer,provider_subject AS subject FROM user_identities')).rows });
mappings: (await pool!.query('SELECT user_id AS "userId",provider AS issuer,provider_subject AS subject FROM user_identities')).rows,
workloads: (await pool!.query('SELECT id,issuer,subject FROM identity_workloads')).rows });
const mapping = { userId: 'preserved', issuer: 'https://identity.example.test', subject: 'first' };
const initial = planIdentityMappings(await inventory(), [mapping]);
await pool.query("UPDATE users SET status='disabled'");
await expect(applyIdentityMappings(database, { requested: [], inventoryDigest: initial.inventoryDigest })).rejects.toThrow('inventory changed');
await expect(applyIdentityMappings(database, { requested: [], ...initial })).rejects.toThrow('inventory or request changed');
await pool.query("UPDATE users SET status='active'");
await applyIdentityMappings(database, { requested: [mapping], inventoryDigest: initial.inventoryDigest });
await expect(applyIdentityMappings(database, { requested: [{ ...mapping, subject: 'changed' }], ...initial })).rejects.toThrow('inventory or request changed');
await pool.query(`INSERT INTO identity_workloads(id,issuer,subject,client_id,display_name,status) VALUES ('service',$1,$2,'client','Service','revoked')`, [mapping.issuer,mapping.subject]);
await expect(applyIdentityMappings(database, { requested: [mapping], ...initial })).rejects.toThrow('workload identity');
await pool.query("DELETE FROM identity_workloads WHERE id='service'");
await applyIdentityMappings(database, { requested: [mapping], ...initial });
const replay = planIdentityMappings(await inventory(), [mapping]);
expect((await applyIdentityMappings(database, { requested: [mapping], inventoryDigest: replay.inventoryDigest })).operations[0].action).toBe('noop');
expect((await applyIdentityMappings(database, { requested: [mapping], ...replay })).operations[0].action).toBe('noop');
await pool.query("ALTER TABLE user_identities ADD CONSTRAINT reject_test_subject CHECK (provider_subject <> 'z-rejected')");
const requested = [{ ...mapping, subject: 'a-inserted' }, { ...mapping, subject: 'z-rejected' }];
const batch = planIdentityMappings(await inventory(), requested);
await expect(applyIdentityMappings(database, { requested, inventoryDigest: batch.inventoryDigest })).rejects.toThrow();
await expect(applyIdentityMappings(database, { requested, ...batch })).rejects.toThrow();
expect((await inventory()).mappings).toEqual([mapping]);
expect((await inventory()).users).toEqual([{ id: 'preserved', status: 'active' }]);
} finally {
Expand Down
Loading