${escapeHtml(setup.secret)}`:label('alreadyShown','p'))+
+ form('/register/totp',ids.formCsrf(token,'registration'),field('otp','otp',{autocomplete:'one-time-code',pattern:'[0-9]{6}',maxlength:6}),'verify'));
+ }
+ if(pathname==='/register/totp'&&request.method==='POST') {
+ submit('registration',['otp']);store.limit(`enroll:${state.account_id}`,10,900);
+ await ids.verifyEnrollment(token,body.get('otp'));redirect(response,'/register/recovery-codes');return true;
+ }
+ if(pathname==='/register/recovery-codes'&&request.method==='GET') {
+ const codes=state.recovery_available?ids.takeRecoveryCodes(token):null;
+ if(!['provisioning','active'].includes(state.status))throw new BoundaryError(400,'REGISTRATION_UNAVAILABLE');
+ return show('recoveryCodes',(codes?label('recoveryNote','p')+`${escapeHtml(state.status)}
Do not enter your ChatGPT password here.
${abort}`, params.redirect_uri); + ${abort}`, params.redirect_uri,config.identity_mode==='multi_account_v1'); } if (prompt.name === "consent") { const scopes = String(params.scope || "").split(" "); + if(enhanced) { + const account=accounts.account(details.session?.accountId); + if(!account||!accounts.eligible(account.subject))throw new BoundaryError(403,'ACCOUNT_DISABLED'); + const keys={openid:'scopeIdentity',offline_access:'scopeOffline','memory:read':'scopeMemory','project:read':'scopeProject'}; + return page(response,` + ${label('oauthClient','p')}${escapeHtml(params.client_id)}This is read-only access to this owner's authorized data, not access restricted to one project. No memory writes, task switching or Resume confirmation.
- ${abort}`, params.redirect_uri); + ${abort}`, params.redirect_uri,config.identity_mode==='multi_account_v1'); } throw new BoundaryError(400, "UNSUPPORTED_INTERACTION"); } @@ -53,10 +80,14 @@ export async function interactionRequest(request, response, { provider, store, a return provider.interactionFinished(request, response, { error: "access_denied", error_description: "The user denied authorization" }, { mergeWithLastSubmission: false }); } if (match[2] === "login" && prompt.name === "login") { - store.limit("login:owner", config.limits.login_attempts_per_account_per_15min, 900); - store.limit(`login:peer:${request.socket.remoteAddress}`, config.limits.login_attempts_per_account_per_15min, 900); + store.limit(`login:account:${String(body.get('username')||'').toLowerCase()}`, config.limits.login_attempts_per_account_per_15min, 900); + store.limit(`login:peer:${request.socket.remoteAddress}`, config.limits.login_attempts_per_account_per_15min*10, 900); const subject = await accounts.authenticate(body.get("username"), body.get("password"), body.get("otp")); if (!subject) throw new BoundaryError(401, "LOGIN_FAILED"); + // A different authenticated principal must start a fresh interaction. Do not + // let the provider's automatic account-switch form reuse the old consent. + if(enhanced && await invalidateBrowserAuthorization(request,response,provider,{nextSubject:subject})) + throw new BoundaryError(409,'AUTHORIZATION_RESTART_REQUIRED'); return provider.interactionFinished(request, response, { login: { accountId: subject, acr: "urn:mnemuron:password-totp", amr: ["pwd", "otp"], ts: seconds() } }, { mergeWithLastSubmission: false }); @@ -65,6 +96,7 @@ export async function interactionRequest(request, response, { provider, store, a const subject = details.session?.accountId; if (!accounts.eligible(subject)) throw new BoundaryError(403, "ACCOUNT_DISABLED"); let grant = details.grantId ? await provider.Grant.find(details.grantId) : undefined; + if(grant && (grant.accountId!==subject || grant.clientId!==params.client_id)) throw new BoundaryError(403,'INTERACTION_MISMATCH'); grant ||= new provider.Grant({ accountId: subject, clientId: params.client_id }); if (prompt.details.missingOIDCScope) grant.addOIDCScope(prompt.details.missingOIDCScope.join(" ")); for (const [resource, scopes] of Object.entries(prompt.details.missingResourceScopes || {})) { diff --git a/services/oauth/src/private-ingress.mjs b/services/oauth/src/private-ingress.mjs new file mode 100644 index 0000000..d409d20 --- /dev/null +++ b/services/oauth/src/private-ingress.mjs @@ -0,0 +1,67 @@ +import tls from 'node:tls'; +import net from 'node:net'; +import {pathToFileURL} from 'node:url'; +import {readPrivate} from '../../../shared/oauth-common.mjs'; + +const privateIPv4=value=>net.isIPv4(value)&&(/^(10\.|192\.168\.)/.test(value)||/^172\.(1[6-9]|2\d|3[01])\./.test(value)); +const integer=(value,min,max)=>Number.isInteger(value)&&value>=min&&value<=max; +export function validateIngressConfig(input,{isolated=false}={}){ + const c={connection_limit:128,handshake_timeout_ms:5000,connect_timeout_ms:3000, + idle_timeout_ms:60000,shutdown_timeout_ms:5000,...input}; + const address=value=>privateIPv4(value)||(isolated&&/^127\./.test(value)&&net.isIPv4(value)); + if(!address(c.listen_host)||!address(c.allowed_peer)||!integer(c.listen_port,isolated?0:1024,65535) + ||!integer(c.upstream_port,1024,65535)||c.upstream_host!==undefined + ||!/^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z][a-z0-9-]{0,62}$/.test(c.server_name||'') + ||!/^[a-f0-9]{64}$/.test(c.client_fingerprint_sha256||'') + ||!integer(c.connection_limit,1,1024) + ||!integer(c.handshake_timeout_ms,100,10000)||!integer(c.connect_timeout_ms,100,10000) + ||!integer(c.idle_timeout_ms,1000,120000)||!integer(c.shutdown_timeout_ms,100,10000)) + throw new Error('PRIVATE_INGRESS_CONFIG_INVALID'); + return c; +} + +// Transport only: the loopback authorization server still owns Host/Origin, +// cookies, MFA, CSRF and OAuth policy. No identity headers are manufactured here. +export function createPrivateIngress(input,{isolated=false}={}){ + const config=validateIngressConfig(input,{isolated}),connections=new Set(); + const server=tls.createServer({key:readPrivate(config.key_file),cert:readPrivate(config.cert_file), + ca:readPrivate(config.ca_file),minVersion:'TLSv1.3',maxVersion:'TLSv1.3', + requestCert:true,rejectUnauthorized:true,handshakeTimeout:config.handshake_timeout_ms, + ALPNProtocols:['http/1.1']},socket=>{ + const fingerprint=socket.getPeerCertificate().fingerprint256?.replaceAll(':','').toLowerCase(); + if(!socket.authorized||socket.remoteAddress!==config.allowed_peer||socket.servername!==config.server_name + ||fingerprint!==config.client_fingerprint_sha256){socket.destroy();return;} + socket.pause(); + const upstream=net.createConnection({host:'127.0.0.1',port:config.upstream_port}); + const timer=setTimeout(()=>{socket.destroy();upstream.destroy();},config.connect_timeout_ms); + timer.unref(); + socket.setTimeout(config.idle_timeout_ms,()=>socket.destroy()); + upstream.setTimeout(config.idle_timeout_ms,()=>upstream.destroy()); + socket.on('error',()=>upstream.destroy());upstream.on('error',()=>socket.destroy()); + socket.on('close',()=>{clearTimeout(timer);upstream.destroy();}); + upstream.on('close',()=>{clearTimeout(timer);socket.destroy();}); + upstream.once('connect',()=>{clearTimeout(timer);socket.pipe(upstream);upstream.pipe(socket);socket.resume();}); + }); + server.on('connection',socket=>{ + if(socket.remoteAddress!==config.allowed_peer||connections.size>=config.connection_limit){socket.destroy();return;} + connections.add(socket);socket.once('close',()=>connections.delete(socket)); + }); + // Do not log peer certificates, request bytes, cookies or TLS error objects. + server.on('tlsClientError',()=>{}); + let closing; + const close=()=>closing??=new Promise(resolve=>{ + if(!server.listening){for(const socket of connections)socket.destroy();resolve();return;} + const timeout=setTimeout(()=>{for(const socket of connections)socket.destroy();},config.shutdown_timeout_ms); + timeout.unref();server.close(()=>{clearTimeout(timeout);resolve();}); + }); + return {server,config,close}; +} + +if(process.argv[1]&&import.meta.url===pathToFileURL(process.argv[1]).href){ + try{ + const ingress=createPrivateIngress(readPrivate(process.argv[2],{json:true})); + ingress.server.on('error',()=>{process.stderr.write('PRIVATE_INGRESS_UNAVAILABLE\n');process.exitCode=1;ingress.close();}); + ingress.server.listen(ingress.config.listen_port,ingress.config.listen_host); + for(const signal of ['SIGINT','SIGTERM'])process.once(signal,()=>{ingress.close();}); + }catch{process.stderr.write('PRIVATE_INGRESS_START_FAILED\n');process.exitCode=1;} +} diff --git a/services/oauth/src/process-lease.mjs b/services/oauth/src/process-lease.mjs new file mode 100644 index 0000000..666f12c --- /dev/null +++ b/services/oauth/src/process-lease.mjs @@ -0,0 +1,26 @@ +import fs from 'node:fs'; +import path from 'node:path'; +import {privateDirectory} from '../../../shared/oauth-common.mjs'; + +// The migration command and server must acquire the same exclusive lease. +export function acquireAuthorizationLease(file) { + const lease=`${file}.process-lock`; + privateDirectory(path.dirname(file),{create:true}); + for(let attempt=0;attempt<2;attempt++) { + try { + const fd=fs.openSync(lease,'wx',0o600); + fs.writeFileSync(fd,String(process.pid));const inode=fs.fstatSync(fd).ino;fs.closeSync(fd); + return ()=>{if(fs.existsSync(lease)&&fs.lstatSync(lease).ino===inode)fs.unlinkSync(lease);}; + } catch(error) { + if(error.code!=='EEXIST')throw error; + const stat=fs.lstatSync(lease); + if(!stat.isFile()||stat.isSymbolicLink()||(stat.mode&0o077)!==0||process.getuid&&stat.uid!==process.getuid())throw new Error('Invalid authorization process lock'); + const pid=Number(fs.readFileSync(lease,'utf8')); + if(!Number.isSafeInteger(pid)||pid<=0)throw new Error('Invalid authorization process lock'); + try{process.kill(pid,0);throw new Error('Only one authorization process may own this database; stop the server before migration');} + catch(probe){if(probe.code!=='ESRCH')throw probe;} + if(fs.lstatSync(lease).ino===stat.ino)fs.unlinkSync(lease); + } + } + throw new Error('Authorization process lock unavailable'); +} diff --git a/services/oauth/src/provider.mjs b/services/oauth/src/provider.mjs index 950492e..f4bc46a 100644 --- a/services/oauth/src/provider.mjs +++ b/services/oauth/src/provider.mjs @@ -53,7 +53,10 @@ export function makeProvider(config, secrets, store, accounts) { }, }, rotateRefreshToken: true, revokeGrantPolicy: () => true, - extraTokenClaims: (_ctx, token) => token.kind === "AccessToken" ? { token_kind: "access_token" } : undefined, + extraTokenClaims: (_ctx, token) => token.kind === "AccessToken" ? { token_kind: "access_token", + ...(accounts.principal?{account_id:accounts.principal(token.accountId).account_id, + security_version:accounts.principal(token.accountId).security_version}:{}), + } : undefined, findAccount: async (_ctx, subject) => accounts.eligible(subject) ? { accountId: subject, claims: async () => ({ sub: subject }) } : undefined, interactions: { url: (_ctx, interaction) => `/interaction/${interaction.uid}` }, diff --git a/services/oauth/src/provisioning.mjs b/services/oauth/src/provisioning.mjs new file mode 100644 index 0000000..ebd0bb2 --- /dev/null +++ b/services/oauth/src/provisioning.mjs @@ -0,0 +1,62 @@ +import fs from 'node:fs'; +import path from 'node:path'; +import {randomUUID} from 'node:crypto'; +import {randomSecret,secretHash,readPrivate,writePrivate,requireConfig,CORE_SCOPES} from '../../../shared/oauth-common.mjs'; +import {storageDoctor} from '../../../server/lib/storage-policy.mjs'; + +export function provisionIdentities(identities,core,{credentialDirectory,identityMapFile,afterCore=()=>{}}) { + storageDoctor({credential_directory:credentialDirectory,identity_map:identityMapFile}); + const db=identities.db; + const operations=db.prepare("SELECT * FROM identity_operations WHERE kind LIKE 'provision:%' AND state NOT IN ('completed','superseded') ORDER BY created,operation_id").all(); + const completed=[]; + for(const operation of operations) { + const a=identities.byId(operation.account_id); + if(a&&operation.kind!==`provision:${a.security_version}`){db.prepare("UPDATE identity_operations SET state='superseded' WHERE operation_id=?").run(operation.operation_id);continue;} + if(!a || !a.mfa_verified || !['provisioning','active'].includes(a.status))continue; + const prepared=identities.store.transaction(()=>{ + const row=db.prepare('SELECT * FROM identity_operations WHERE operation_id=?').get(operation.operation_id); + if(row.payload_cipher)return identities.unseal(row.payload_cipher,a.account_id,'provision'); + const bindings=['web','console'].map(purpose=>({purpose,credential_id:randomUUID(),api_key:`mnm_${randomSecret()}`, + user_id:a.user_id,agent_instance_id:`${purpose}-${a.account_id}`,agent_id:purpose==='web'?'chatgpt-web':'mnemuron-console', + scopes:purpose==='web'?[...CORE_SCOPES]:['memory:read','resume:read','console:read'], + credential_file:path.join(credentialDirectory,`${a.account_id}-${purpose}-v${a.security_version}.key`)})); + db.prepare("UPDATE identity_operations SET state='prepared',payload_cipher=? WHERE operation_id=?") + .run(identities.seal(bindings,a.account_id,'provision'),operation.operation_id); + return bindings; + }); + try { + // The durable operation is committed BEFORE Core. An uncertain completion reuses the same key/id. + core.memoryTransaction(()=>{ + for(const b of prepared) { + const existing=core.db.prepare('SELECT * FROM credentials WHERE credential_id=? OR key_hash=?').get(b.credential_id,secretHash(b.api_key)); + if(existing) { + requireConfig(existing.credential_id===b.credential_id&&existing.user_id===a.user_id&&existing.agent_instance_id===b.agent_instance_id + &&existing.agent_id===b.agent_id&&existing.key_hash===secretHash(b.api_key)&&!existing.revoked_at&&existing.scopes_json===JSON.stringify(b.scopes),'core provisioning conflict'); + } else core.db.prepare(`INSERT INTO credentials(credential_id,label,user_id,device_id,agent_id,agent_instance_id,key_hash,scopes_json,created_at) + VALUES(?,?,?,?,?,?,?,?,?)`).run(b.credential_id,'Account-bound read-only connection',a.user_id,'identity-service',b.agent_id,b.agent_instance_id,secretHash(b.api_key),JSON.stringify(b.scopes),new Date().toISOString()); + } + }); + afterCore(operation); + for(const b of prepared) { + const auth=core.authenticate(b.api_key); + requireConfig(auth.user_id===a.user_id&&auth.credential_id===b.credential_id&&auth.agent_instance_id===b.agent_instance_id,'authoritative core identity'); + if(fs.existsSync(b.credential_file))requireConfig(readPrivate(b.credential_file)===b.api_key,'credential file collision'); + else writePrivate(b.credential_file,b.api_key); + } + identities.finishProvision(a.account_id,prepared,operation.operation_id); + completed.push(a.account_id); + } catch(error) { + db.prepare("UPDATE identity_operations SET last_error='PROVISIONING_INCOMPLETE' WHERE operation_id=?").run(operation.operation_id); + throw error; + } + } + // Derived publication only: OAuth introspection remains authoritative for current activity/version. + identities.store.transaction(()=>{ + const mappings=db.prepare(`SELECT a.*,b.credential_id,b.agent_instance_id,b.credential_file FROM identity_accounts a + JOIN identity_bindings b ON a.account_id=b.account_id AND b.purpose='web' AND b.checked=1 WHERE a.binding_ready=1`).all().map(a=>({ + issuer:a.issuer,subject:a.subject,account_id:a.account_id,mnemuron_user_id:a.user_id,security_version:a.security_version, + enabled:a.status==='active',agent_instance_id:a.agent_instance_id,credential_id:a.credential_id,credential_file:a.credential_file})); + writePrivate(identityMapFile,{schema_version:'multi-account-identity-v1',unknown_subject_policy:'deny',mappings},{replace:fs.existsSync(identityMapFile)}); + }); + return {completed:completed.length,pending:db.prepare("SELECT COUNT(*) n FROM identity_operations WHERE kind LIKE 'provision:%' AND state NOT IN ('completed','superseded')").get().n}; +} diff --git a/services/oauth/src/recovery.mjs b/services/oauth/src/recovery.mjs new file mode 100644 index 0000000..9dce10b --- /dev/null +++ b/services/oauth/src/recovery.mjs @@ -0,0 +1,102 @@ +import {scrypt,randomUUID} from 'node:crypto'; +import {promisify} from 'node:util'; +import {generateSecret,generateURI,verify} from 'otplib'; +import {BoundaryError,secretHash,equalSecret,seconds} from '../../../shared/oauth-common.mjs'; +import {usernameKey,passwordRecord} from './identity-repository.mjs'; +const derive=promisify(scrypt); +const fail=()=>{throw new BoundaryError(400,'RECOVERY_UNAVAILABLE');}; + +// No constructor default approves a proof combination. HTTP/CLI remain blocked +// until an operator provides an independently approved policy integration. +export class RecoveryService { + constructor(identities,{policy=null}={}) {this.ids=identities;this.db=identities.db;this.policy=policy;} + proofs(action) { + const required=this.policy?.[action]; + if(!['password','totp'].includes(action)||!Array.isArray(required)||required.length!==2||!required.includes('recovery_code') + ||!required.includes(action==='password'?'totp':'password'))throw new BoundaryError(403,'BLOCKED_POLICY'); + return required; + } + async begin({username,action,password,otp,recoveryCode}) { + const required=this.proofs(action),key=usernameKey(username); + this.ids.store.limit(`recovery:${key}`,5,900); + const a=this.db.prepare('SELECT * FROM identity_accounts WHERE username_key=?').get(key); + if(!a||!this.ids.eligible(a.subject)||typeof recoveryCode!=='string'||recoveryCode.length>256)fail(); + const digest=secretHash(recoveryCode),hashes=JSON.parse(a.recovery_hashes); + if(!hashes.some(h=>equalSecret(h,digest)))fail(); + let epoch; + if(required.includes('password')) { + const p=JSON.parse(a.password_json),supplied=typeof password==='string'&&password.length<=1024?password:''; + const actual=(await derive(supplied,p.salt,64,{N:p.N,r:p.r,p:p.p,maxmem:128*1024*1024})).toString('base64url'); + if(!equalSecret(actual,p.hash))fail(); + } + if(required.includes('totp')) { + if(!/^\d{6}$/.test(otp||''))fail(); + const result=await verify({secret:this.ids.unseal(a.mfa_cipher,a.account_id,'totp'),token:otp,epochTolerance:30}); + if(!result.valid)fail();epoch=result.epoch; + } + return this.ids.store.transaction(()=>{ + const current=this.ids.byId(a.account_id); + if(!this.ids.eligible(a.subject)||current.security_version!==a.security_version||!JSON.parse(current.recovery_hashes).includes(digest))fail(); + if(epoch!==undefined&&!this.ids.store.consumeStep(a.subject,epoch))fail(); + this.db.prepare("UPDATE identity_accounts SET status='recovery_pending',security_version=security_version+1,recovery_hashes=? WHERE account_id=?") + .run(JSON.stringify(hashes.filter(h=>h!==digest)),a.account_id); + const s=this.ids.newSession('recovery',{accountId:a.account_id,ttl:600}); + this.db.prepare('INSERT INTO identity_recovery_claims VALUES(?,?,?)').run(a.account_id,digest,secretHash(s.token)); + const id=randomUUID(),payload={session_digest:secretHash(s.token),action,...(action==='totp'?{totp_secret:generateSecret()}:{}),prior_security_version:a.security_version}; + this.db.prepare('INSERT INTO identity_operations(operation_id,account_id,kind,state,payload_cipher,created) VALUES(?,?,?,?,?,?)') + .run(id,a.account_id,`recovery:${id}`,'proof_verified',this.ids.seal(payload,a.account_id,'recovery'),seconds()); + this.ids.audit(a.account_id,'recovery.proof_verified');return {token:s.token,csrf:s.csrf,restricted:true,action}; + }); + } + operation(token) { + const s=this.ids.session(token,'recovery'); + const rows=this.db.prepare("SELECT * FROM identity_operations WHERE account_id=? AND kind LIKE 'recovery:%'").all(s.account_id); + const op=rows.map(row=>({...row,payload:this.ids.unseal(row.payload_cipher,s.account_id,'recovery')})).find(row=>row.payload.session_digest===s.digest); + if(!op)fail();return op; + } + enrollment(token) { + const op=this.operation(token);if(op.payload.action!=='totp'||op.state!=='proof_verified')fail(); + const a=this.ids.byId(op.account_id);return {uri:generateURI({issuer:'Mnemuron',label:a.username,secret:op.payload.totp_secret}),secret:op.payload.totp_secret}; + } + async complete(token,{password,otp},{revokeCore}={}) { + let op=this.operation(token),a=this.ids.byId(op.account_id);this.proofs(op.payload.action); + if(op.state==='completed')return {status:a.status,login_required:true}; + if(typeof revokeCore!=='function')throw new BoundaryError(503,'REVOCATION_DEPENDENCY_REQUIRED'); + if(op.state==='proof_verified') { + let record,epoch; + if(op.payload.action==='password')record=await passwordRecord(password); + else { + if(!/^\d{6}$/.test(otp||''))fail();const v=await verify({secret:op.payload.totp_secret,token:otp,epochTolerance:30}); + if(!v.valid)fail();epoch=v.epoch; + } + this.ids.store.transaction(()=>{ + const current=this.operation(token);if(current.state!=='proof_verified')fail(); + if(record)this.db.prepare('UPDATE identity_accounts SET password_json=? WHERE account_id=?').run(JSON.stringify(record),a.account_id); + else { + if(!this.ids.store.consumeStep(a.subject,epoch))fail(); + this.db.prepare('UPDATE identity_accounts SET mfa_cipher=?,mfa_verified=1 WHERE account_id=?').run(this.ids.seal(op.payload.totp_secret,a.account_id,'totp'),a.account_id); + } + this.db.prepare("UPDATE identity_operations SET state='revocation_pending' WHERE operation_id=?").run(op.operation_id); + this.ids.audit(a.account_id,'recovery.credentials_prepared'); + }); + } + // A durable paused account invalidates introspection and console sessions even + // when the independent Core revocation is interrupted. No success is claimed. + try { + this.ids.store.revoke({subject:a.subject}); + const bindings=this.db.prepare('SELECT credential_id FROM identity_bindings WHERE account_id=?').all(a.account_id); + if(await revokeCore({user_id:a.user_id,credential_ids:bindings.map(b=>b.credential_id)})!==true)throw new Error('Core revocation unverified'); + }catch { + this.db.prepare("UPDATE identity_operations SET last_error='REVOCATION_INCOMPLETE' WHERE operation_id=?").run(op.operation_id); + throw new BoundaryError(503,'REVOCATION_INCOMPLETE'); + } + this.ids.store.transaction(()=>{ + this.operation(token); + this.db.prepare("UPDATE identity_accounts SET status='provisioning',binding_ready=0 WHERE account_id=? AND status='recovery_pending'").run(a.account_id); + this.ids.queueProvision(a.account_id); + this.db.prepare("UPDATE identity_operations SET state='completed',last_error=NULL WHERE operation_id=?").run(op.operation_id); + this.ids.audit(a.account_id,'recovery.completed'); + }); + return {status:'provisioning',login_required:true}; + } +} diff --git a/services/oauth/src/server.mjs b/services/oauth/src/server.mjs index fd30444..a10e0ce 100644 --- a/services/oauth/src/server.mjs +++ b/services/oauth/src/server.mjs @@ -1,39 +1,20 @@ import http from "node:http"; -import fs from "node:fs"; import { pathToFileURL } from "node:url"; import { randomUUID } from "node:crypto"; import { loadAuthConfig, validateAuthConfig, loadAuthSecrets } from "./config.mjs"; import { AuthStore } from "./sqlite-adapter.mjs"; import { Accounts } from "./accounts.mjs"; +import { IdentityRepository } from './identity-repository.mjs'; +import {consoleRequest} from './console.mjs'; +import {ConsoleCore} from './console-core.mjs'; +import {sendPage,label} from '../../../web/console/render.mjs'; import { makeProvider } from "./provider.mjs"; import { interactionRequest } from "./interactions.mjs"; import { BoundaryError, SerialGate, WindowLimit, OAUTH_SCOPES, parseForm, readBody, - requestBoundary, sendJson, equalSecret, privateDirectory } from "../../../shared/oauth-common.mjs"; -import path from "node:path"; - -function acquireLease(file) { - const lease = `${file}.process-lock`; - privateDirectory(path.dirname(file), { create: true }); - for (let attempt = 0; attempt < 2; attempt++) { - try { - const fd = fs.openSync(lease, "wx", 0o600); - fs.writeFileSync(fd, String(process.pid)); - const inode = fs.fstatSync(fd).ino; - fs.closeSync(fd); - return () => { if (fs.existsSync(lease) && fs.lstatSync(lease).ino === inode) fs.unlinkSync(lease); }; - } catch (error) { - if (error.code !== "EEXIST") throw error; - const stat = fs.lstatSync(lease); - if (!stat.isFile() || stat.isSymbolicLink()) throw new Error("Invalid authorization process lock"); - const pid = Number(fs.readFileSync(lease, "utf8")); - if (!Number.isInteger(pid) || pid <= 0) throw new Error("Invalid authorization process lock"); - try { process.kill(pid, 0); throw new Error("Only one authorization server may own this database"); } - catch (probe) { if (probe.code !== "ESRCH") throw probe; } - if (fs.lstatSync(lease).ino === stat.ino) fs.unlinkSync(lease); - } - } - throw new Error("Authorization process lock unavailable"); -} + requestBoundary, sendJson, equalSecret } from "../../../shared/oauth-common.mjs"; +import {acquireAuthorizationLease} from './process-lease.mjs'; +import {storageDoctor} from '../../../server/lib/storage-policy.mjs'; +import {invalidateBrowserAuthorization} from './browser-session.mjs'; function bootstrapMetadata(config) { return { issuer: config.issuer, authorization_endpoint: `${config.issuer}/authorize`, @@ -61,15 +42,22 @@ export function createAuthorizationServer(input, { isolated = false, logger = () let store, accounts, provider, secrets, release; try { if (config.mode === "oauth") { + if(config.identity_mode==='multi_account_v1') storageDoctor({ + identity_database:config.database_file,identity_key:config.identity.encryption_key_file, + }); secrets = loadAuthSecrets(config); - release = acquireLease(config.database_file); - store = new AuthStore(config.database_file); - accounts = new Accounts(config.accounts_file, store); - accounts.read(); + release = acquireAuthorizationLease(config.database_file); + store = new AuthStore(config.database_file,{identity:config.identity_mode==='multi_account_v1'}); + accounts = config.identity_mode==='multi_account_v1' ? new IdentityRepository(store,{ + keyFile:config.identity.encryption_key_file,issuer:config.issuer,batchLimit:config.identity.invitation_batch_limit, + sessionTtl:config.identity.console_session_ttl_seconds}) : new Accounts(config.accounts_file, store); + if(config.identity_mode==='legacy_owner') accounts.read(); + else store.identity=accounts; provider = makeProvider(config, secrets, store, accounts); } } catch (error) { store?.close(); release?.(); throw error; } const gate = new SerialGate(); + const consoleGate = new SerialGate(); const limits = new WindowLimit(); const callback = provider?.callback(); const origin = new URL(config.issuer); @@ -89,7 +77,8 @@ export function createAuthorizationServer(input, { isolated = false, logger = () const url = requestBoundary(request, origin, { isolated }); limits.take(`peer:${request.socket.remoteAddress}`, 1500); if (request.method === "GET" && ["/livez", "/readyz"].includes(url.pathname)) { - const ready = config.mode === "oauth" && store.ready({ writeProbe: url.pathname === "/readyz" }) && accounts.read().enabled; + const ready = config.mode === "oauth" && store.ready({ writeProbe: url.pathname === "/readyz" }) + && (config.identity_mode==='multi_account_v1' || accounts.read().enabled); return sendJson(response, url.pathname === "/livez" || ready ? 200 : 503, { service: "mnemuron-oauth", mode: config.mode, ready, production_ready: false }); } @@ -107,6 +96,10 @@ export function createAuthorizationServer(input, { isolated = false, logger = () request.url = "/.well-known/openid-configuration"; return callback(request, response); } + const handleConsole=()=>consoleRequest(request,response,{config,accounts,store,url, + invalidateAuthorization:nextSubject=>invalidateBrowserAuthorization(request,response,provider,{nextSubject}),coreFor:subject=>new ConsoleCore(config.identity?.core, + accounts.principal(subject),accounts.bindings(subject).find(b=>b.purpose==='console'))}); + if(await (request.method==='POST'&&/^\/(register|login)(\/|$)/.test(url.pathname)?consoleGate.run(handleConsole):handleConsole()))return; if (url.pathname.startsWith("/interaction/")) { return await gate.run(() => interactionRequest(request, response, { provider, store, accounts, config, url })); } @@ -163,6 +156,16 @@ export function createAuthorizationServer(input, { isolated = false, logger = () return await gate.run(() => callback(request, response)); } catch (error) { errorCode = error instanceof BoundaryError ? error.code : "AUTH_DEPENDENCY_UNAVAILABLE"; + const route=request.url.split('?')[0],interaction=route.match(/^\/interaction\/([A-Za-z0-9_-]{1,128})(?:\/(login|confirm|abort))?$/); + if(!response.headersSent&&config.identity_mode==='multi_account_v1'&&request.headers.accept?.includes('text/html') + && (/^\/(register|login|recover)(\/|$)/.test(route)||interaction)) { + const status=error instanceof BoundaryError?error.status:503; + const restart=['AUTHORIZATION_RESTART_REQUIRED','INTERACTION_EXPIRED','INTERACTION_MISMATCH'].includes(errorCode); + const message=restart?'restartAuthorization':errorCode==='LOGIN_FAILED'?'loginFailed':status===429?'rateLimited':status>=500?'unavailable':errorCode==='BLOCKED_POLICY'?'blockedNote':'pendingStep'; + const back=interaction&&!restart?`/interaction/${interaction[1]}`:route.startsWith('/register')?'/register':'/login'; + const navigation=interaction&&restart?label('oauthRestartHelp','p'):`${label(interaction?'oauthRetry':'back')}`; + sendPage(response,{title:'error',auth:true,authPurpose:interaction?'oauth':'console',body:`