diff --git a/changelog/1.0.0-rc.5/fs-authority-native.md b/changelog/1.0.0-rc.5/fs-authority-native.md new file mode 100644 index 00000000..c38ca4f5 --- /dev/null +++ b/changelog/1.0.0-rc.5/fs-authority-native.md @@ -0,0 +1,3 @@ +# @oneworks/fs-authority-native 1.0.0-rc.5 + +- Fix a macOS first-startup race that could make an asset authority request fail before the local filesystem authority is ready. diff --git a/packages/fs-authority-native/__tests__/broker-ownership.test.cjs b/packages/fs-authority-native/__tests__/broker-ownership.test.cjs index e7087ba8..60dd0de8 100644 --- a/packages/fs-authority-native/__tests__/broker-ownership.test.cjs +++ b/packages/fs-authority-native/__tests__/broker-ownership.test.cjs @@ -3,6 +3,7 @@ process.env.NODE_ENV = 'test' const assert = require('node:assert/strict') +const { once } = require('node:events') const { mkdtempSync, rmSync } = require('node:fs') const { tmpdir } = require('node:os') const { join } = require('node:path') @@ -13,6 +14,11 @@ const { startFilesystemAuthorityBroker } = require('../testing.cjs') +const getConnectionCount = server => + new Promise((resolve, reject) => { + server.getConnections((error, count) => error == null ? resolve(count) : reject(error)) + }) + test('a delayed startup loser cannot recover or clear the endpoint winner claim', async () => { const root = mkdtempSync(join(tmpdir(), 'ow-authority-owner-')) const workspace = mkdtempSync(join(root, 'workspace-')) @@ -79,3 +85,137 @@ test('a delayed startup loser cannot recover or clear the endpoint winner claim' rmSync(root, { force: true, recursive: true }) } }) + +test('holds a trusted local connection until broker recovery is ready', async () => { + const root = mkdtempSync(join(tmpdir(), 'ow-authority-readiness-')) + const workspace = mkdtempSync(join(root, 'workspace-')) + const prepared = prepareFilesystemAuthorityTestControlRoot(join(root, 'control')) + let enterRecover + let releaseRecover + const entered = new Promise(resolve => { + enterRecover = resolve + }) + const released = new Promise(resolve => { + releaseRecover = resolve + }) + const brokerPromise = startFilesystemAuthorityBroker({ + beforeRecover: context => { + enterRecover(context) + return released + }, + controlRoot: prepared.controlRoot, + secret: prepared.secret + }) + const { server } = await entered + const connected = once(server, 'connection') + const authorityPromise = openFilesystemAuthorityForTest(workspace, { + autoStart: false, + controlRoot: prepared.controlRoot, + secret: prepared.secret, + timeoutMs: 1000 + }) + await connected + + releaseRecover() + const [broker, authority] = await Promise.all([brokerPromise, authorityPromise]) + try { + assert.equal(await authority.claimMutation('broker-readiness', 'slow-start'), 1) + } finally { + authority.close() + await broker.close() + rmSync(root, { force: true, recursive: true }) + } +}) + +test('cleans up a connection that times out while broker recovery never becomes ready', async () => { + const root = mkdtempSync(join(tmpdir(), 'ow-authority-timeout-')) + const workspace = mkdtempSync(join(root, 'workspace-')) + const prepared = prepareFilesystemAuthorityTestControlRoot(join(root, 'control')) + let enterRecover + let releaseRecover + const entered = new Promise(resolve => { + enterRecover = resolve + }) + const released = new Promise(resolve => { + releaseRecover = resolve + }) + const brokerPromise = startFilesystemAuthorityBroker({ + beforeRecover: context => { + enterRecover(context) + return released + }, + controlRoot: prepared.controlRoot, + secret: prepared.secret + }) + const { server } = await entered + let socket + const closed = once(server, 'connection').then(([connection]) => { + socket = connection + return once(socket, 'close') + }) + const authorityPromise = openFilesystemAuthorityForTest(workspace, { + autoStart: false, + controlRoot: prepared.controlRoot, + secret: prepared.secret, + timeoutMs: 25 + }) + await assert.rejects( + authorityPromise, + { code: 'asset_filesystem_authority_unavailable', committed: false } + ) + await closed + assert.equal(socket.listenerCount('data'), 0) + + releaseRecover() + const broker = await brokerPromise + try { + assert.equal(await getConnectionCount(broker.server), 0) + } finally { + await broker.close() + rmSync(root, { force: true, recursive: true }) + } +}) + +test('rejects and cleans up pending connections when broker recovery fails terminally', async () => { + const root = mkdtempSync(join(tmpdir(), 'ow-authority-fail-')) + const workspace = mkdtempSync(join(root, 'workspace-')) + const prepared = prepareFilesystemAuthorityTestControlRoot(join(root, 'control')) + let enterRecover + let failRecover + const entered = new Promise(resolve => { + enterRecover = resolve + }) + const failed = new Promise(resolve => { + failRecover = resolve + }) + const brokerPromise = startFilesystemAuthorityBroker({ + beforeRecover: async context => { + enterRecover(context) + await failed + throw new Error('recovery failed') + }, + controlRoot: prepared.controlRoot, + secret: prepared.secret + }) + const { server } = await entered + let socket + const connected = once(server, 'connection').then(([connection]) => { + socket = connection + }) + const authorityPromise = openFilesystemAuthorityForTest(workspace, { + autoStart: false, + controlRoot: prepared.controlRoot, + secret: prepared.secret, + timeoutMs: 1000 + }) + await connected + failRecover() + + await assert.rejects(brokerPromise, /recovery failed/u) + await assert.rejects( + authorityPromise, + { code: 'asset_filesystem_authority_unavailable', committed: false } + ) + assert.equal(socket.listenerCount('data'), 0) + rmSync(root, { force: true, recursive: true }) +}) diff --git a/packages/fs-authority-native/broker-session.cjs b/packages/fs-authority-native/broker-session.cjs index a70dd086..9f7e1dd2 100644 --- a/packages/fs-authority-native/broker-session.cjs +++ b/packages/fs-authority-native/broker-session.cjs @@ -16,9 +16,9 @@ const { createFrameChannel, encodeFrame } = require('./protocol.cjs') const { writeSocket } = require('./transport.cjs') const failure = (code, committed = false, warnings = []) => ({ error: { code, committed, warnings }, ok: false }) const createSession = ( - { allowFaults, binding, claims, controlRoot, database, epoch, secret, socket, workspaceRoot } + { allowFaults, binding, claims, controlRoot, database, epoch, initialBytes, secret, socket, workspaceRoot } ) => { - const channel = createFrameChannel(socket) + const channel = createFrameChannel(socket, initialBytes) const handshake = createBrokerHandshake(epoch, secret) const managedTrees = createManagedTreeHandler({ allowFaults, binding, database, secret }) const state = { authenticated: false, authority: undefined, claim: undefined, socket } diff --git a/packages/fs-authority-native/broker.cjs b/packages/fs-authority-native/broker.cjs index e112d141..5567b49f 100644 --- a/packages/fs-authority-native/broker.cjs +++ b/packages/fs-authority-native/broker.cjs @@ -1,13 +1,21 @@ #!/usr/bin/env node 'use strict' +const { Buffer } = require('node:buffer') const { randomBytes } = require('node:crypto') const { resolve } = require('node:path') const process = require('node:process') const { createSession } = require('./broker-session.cjs') const { openClaimDatabase } = require('./claim-db.cjs') -const { canonicalWorkspace, prepareControlRoot, readOrCreateSecret, secureBrokerEndpoint } = require('./constants.cjs') +const { + canonicalWorkspace, + MAX_FRAME_BYTES, + prepareControlRoot, + readOrCreateSecret, + secureBrokerEndpoint +} = require('./constants.cjs') const { loadBinding } = require('./loader.cjs') const { createBrokerServer, prepareEndpointForListen, verifySocketPeer } = require('./transport.cjs') +const MAX_PENDING_CONNECTIONS = 64 const readTestConfig = () => { if (process.env.NODE_ENV !== 'test') return {} const index = process.argv.indexOf('--test-control-root') @@ -48,15 +56,13 @@ const startBroker = async ( const endpoint = await prepareEndpointForListen(controlRoot) const claims = new Map() const epoch = randomBytes(32).toString('hex') + const pendingSockets = new Map() const sessions = new Set() let ready = false let database - const onConnection = socket => { - if (!ready) { - socket.destroy() - return - } - if (!verifySocketPeer(binding, socket, true)) { + const startSession = (socket, initialBytes, peerVerified = false) => { + if (socket.destroyed) return + if (!peerVerified && !verifySocketPeer(binding, socket, true)) { socket.destroy() return } @@ -68,6 +74,7 @@ const startBroker = async ( controlRoot, database, epoch, + initialBytes, secret, socket, workspaceRoot: canonicalWorkspace @@ -75,21 +82,71 @@ const startBroker = async ( sessions.add(session) session.done = session.run().finally(() => sessions.delete(session)) } + const cleanupPendingSocket = (socket, pending) => { + if (pending.cleaned) return + pending.cleaned = true + socket.off('data', pending.onData) + socket.off('close', pending.onClose) + socket.off('error', pending.onError) + if (pendingSockets.get(socket) === pending) pendingSockets.delete(socket) + pending.chunks.length = 0 + pending.bytes = 0 + } + const destroyPendingSockets = () => { + for (const [socket, pending] of [...pendingSockets]) { + cleanupPendingSocket(socket, pending) + socket.destroy() + } + } + const onConnection = socket => { + if (ready) { + startSession(socket) + return + } + if (!verifySocketPeer(binding, socket, true) || pendingSockets.size >= MAX_PENDING_CONNECTIONS) { + socket.destroy() + return + } + socket.setNoDelay(true) + const pending = { bytes: 0, chunks: [], cleaned: false } + pending.onData = chunk => { + pending.bytes += chunk.length + if (pending.bytes > MAX_FRAME_BYTES + 4) { + socket.destroy() + return + } + pending.chunks.push(chunk) + } + pending.onClose = () => cleanupPendingSocket(socket, pending) + pending.onError = () => cleanupPendingSocket(socket, pending) + socket.on('data', pending.onData) + pendingSockets.set(socket, pending) + socket.once('close', pending.onClose) + socket.once('error', pending.onError) + } let server try { server = createBrokerServer(binding, endpoint, onConnection) await listen(server, endpoint, controlRoot) database = suppliedDatabase ?? openClaimDatabase(controlRoot) - await beforeRecover?.() + await beforeRecover?.({ endpoint, server }) database.recover(epoch) ready = true + for (const [socket, pending] of [...pendingSockets]) { + const initialBytes = Buffer.concat(pending.chunks, pending.bytes) + cleanupPendingSocket(socket, pending) + startSession(socket, initialBytes, true) + } } catch (error) { + ready = false + destroyPendingSockets() await closeServer(server) database?.close() throw error } const close = async () => { const active = [...sessions] + destroyPendingSockets() for (const session of active) session.state.socket?.destroy?.() ready = false await Promise.all([closeServer(server), ...active.map(session => session.done)]) diff --git a/packages/fs-authority-native/protocol.cjs b/packages/fs-authority-native/protocol.cjs index 8b770f2d..a3bca5f7 100644 --- a/packages/fs-authority-native/protocol.cjs +++ b/packages/fs-authority-native/protocol.cjs @@ -10,7 +10,7 @@ const encodeFrame = value => { payload.copy(frame, 4) return frame } -const createFrameChannel = stream => { +const createFrameChannel = (stream, initialBytes) => { const queued = [] const waiting = [] let buffered = Buffer.alloc(0) @@ -30,7 +30,7 @@ const createFrameChannel = stream => { queued.push(value) return true } - stream.on('data', chunk => { + const onData = chunk => { if (buffered.length + chunk.length > MAX_FRAME_BYTES + 4) { fail(new Error('Filesystem authority channel buffer is full')) stream.destroy() @@ -63,9 +63,11 @@ const createFrameChannel = stream => { return } } - }) + } + stream.on('data', onData) stream.once('error', fail) stream.once('close', () => fail(new Error('Filesystem authority connection closed'))) + if (initialBytes?.length > 0) onData(initialBytes) return { next(timeoutMs) { if (queued.length > 0) return Promise.resolve(queued.shift())