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
3 changes: 3 additions & 0 deletions changelog/1.0.0-rc.5/fs-authority-native.md
Original file line number Diff line number Diff line change
@@ -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.
140 changes: 140 additions & 0 deletions packages/fs-authority-native/__tests__/broker-ownership.test.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand All @@ -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-'))
Expand Down Expand Up @@ -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 })
})
4 changes: 2 additions & 2 deletions packages/fs-authority-native/broker-session.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
73 changes: 65 additions & 8 deletions packages/fs-authority-native/broker.cjs
Original file line number Diff line number Diff line change
@@ -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')
Expand Down Expand Up @@ -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
}
Expand All @@ -68,28 +74,79 @@ const startBroker = async (
controlRoot,
database,
epoch,
initialBytes,
secret,
socket,
workspaceRoot: canonicalWorkspace
})
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)])
Expand Down
8 changes: 5 additions & 3 deletions packages/fs-authority-native/protocol.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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()
Expand Down Expand Up @@ -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())
Expand Down
Loading