From 35a16706c70a12ace30064b3e5b19688b78f9712 Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 11:28:22 +0200 Subject: [PATCH 1/7] Fix overlapping CLI watch rebuilds --- .../src/commands/build/index.test.ts | 78 ++++++++ .../src/commands/build/index.ts | 172 +++++++++++++----- .../src/utils/serial-batches.test.ts | 155 ++++++++++++++++ .../src/utils/serial-batches.ts | 65 +++++++ 4 files changed, 420 insertions(+), 50 deletions(-) create mode 100644 packages/@tailwindcss-cli/src/commands/build/index.test.ts create mode 100644 packages/@tailwindcss-cli/src/utils/serial-batches.test.ts create mode 100644 packages/@tailwindcss-cli/src/utils/serial-batches.ts diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts new file mode 100644 index 000000000..dbc0cc42d --- /dev/null +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -0,0 +1,78 @@ +import { expect, it } from 'vitest' +import { serializeBatches } from '../../utils/serial-batches' +import { createWatchers, filterChangedFiles } from './index' + +type WatchEvent = { type: 'create' | 'update' | 'delete'; path: string } +type WatchCallback = (error: Error | null, events: WatchEvent[]) => Promise + +function fakeWatcher() { + let callbacks: WatchCallback[] = [] + return { + callbacks, + watcher: { + async subscribe(_directory: string, callback: WatchCallback) { + callbacks.push(callback) + return { unsubscribe() {} } + }, + }, + } +} + +function nextTask() { + return new Promise((resolve) => setTimeout(resolve, 0)) +} + +it('removes duplicate output and map events from a coalesced batch', () => { + expect( + filterChangedFiles( + ['output.css', 'source.html', 'output.css', 'output.css.map'], + 'output.css', + 'output.css.map', + ), + ).toEqual(['source.html']) +}) + +it('flushes a collected event when shutdown cancels its debounce timer', async () => { + let calls: string[][] = [] + let queue = serializeBatches(async (files) => { + calls.push(files) + }) + let fake = fakeWatcher() + let generation = await createWatchers(['/watch'], async () => {}, queue, fake.watcher) + + await fake.callbacks[0](null, [{ type: 'delete', path: 'last-change' }]) + await generation.cleanup() + await queue.close() + + expect(calls).toEqual([['last-change']]) +}) + +it('waits for an entered watcher callback before shutdown flushes changes', async () => { + let releaseLstat!: () => void + let lstatCanFinish = new Promise((resolve) => (releaseLstat = resolve)) + let calls: string[][] = [] + let queue = serializeBatches(async (files) => { + calls.push(files) + }) + let fake = fakeWatcher() + let generation = await createWatchers( + ['/watch'], + async () => {}, + queue, + fake.watcher, + async () => { + await lstatCanFinish + return { isFile: () => true, isSymbolicLink: () => false } + }, + ) + + let callback = fake.callbacks[0](null, [{ type: 'update', path: 'delayed-change' }]) + let cleanup = generation.cleanup() + await nextTask() + expect(calls).toEqual([]) + + releaseLstat() + await Promise.all([callback, cleanup]) + await queue.close() + expect(calls).toEqual([['delayed-change']]) +}) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.ts b/packages/@tailwindcss-cli/src/commands/build/index.ts index c85eda28c..950692939 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.ts @@ -22,6 +22,7 @@ import { relative, wordWrap, } from '../../utils/renderer' +import { serializeBatches, type SerialBatches } from '../../utils/serial-batches' import { drainStdin, outputFile } from './utils' const css = String.raw @@ -326,6 +327,13 @@ export async function handle(args: Result>) { let [compiler, scanner] = await handleError(() => createCompiler(input, I)) let cleanupWatchers: (() => Promise)[] = [] + let finishInitialBuild!: () => void + let initialBuildFinished = new Promise((resolve) => (finishInitialBuild = resolve)) + let eventBatches: SerialBatches | null = null + let setEventHandler!: (handler: (files: string[]) => Promise) => void + let eventHandler = new Promise<(files: string[]) => Promise>( + (resolve) => (setEventHandler = resolve), + ) // Watch for changes if (args['--watch'] && pollInterval === false) { @@ -333,12 +341,18 @@ export async function handle(args: Result>) { // such that we can present a helpful error message if needed. await handleError(() => loadWatcher()) - cleanupWatchers.push( - await createWatchers(await watchDirectories(scanner), async function handle(files) { + eventBatches = serializeBatches( + async (files) => (await eventHandler)(files), + initialBuildFinished, + (error) => eprintln(formatError(error)), + ) + let initialWatchers = await createWatchers( + await watchDirectories(scanner), + async function handle(files) { try { - // If the only change happened to the output file, then we don't want to - // trigger a rebuild because that will result in an infinite loop. - if (files.length === 1 && files[0] === args['--output']) return + // Ignore our own writes so they don't trigger another rebuild. + files = filterChangedFiles(files, args['--output'], args['--map']) + if (files.length === 0) return using I = new Instrumentation() DEBUG && I.start('[@tailwindcss/cli] (watcher)') @@ -388,7 +402,11 @@ export async function handle(args: Result>) { // Setup new watchers DEBUG && I.start('Setup new watchers') - let newCleanupFunction = await createWatchers(await watchDirectories(scanner), handle) + let newWatchers = await createWatchers( + await watchDirectories(scanner), + handle, + eventBatches!, + ) DEBUG && I.end('Setup new watchers') // Clear old watchers @@ -396,7 +414,7 @@ export async function handle(args: Result>) { await Promise.all(cleanupWatchers.splice(0).map((cleanup) => cleanup())) DEBUG && I.end('Cleanup old watchers') - cleanupWatchers.push(newCleanupFunction) + cleanupWatchers.push(newWatchers.cleanup) // Re-compile the CSS DEBUG && I.start('Build CSS') @@ -481,17 +499,22 @@ export async function handle(args: Result>) { let end = process.hrtime.bigint() if (!args['--silent']) eprintln(`Done in ${formatDuration(end - start)}`) } - }), + }, + eventBatches, ) + setEventHandler(initialWatchers.callback) + cleanupWatchers.push(initialWatchers.cleanup) // Abort the watcher if `stdin` is closed to avoid zombie processes. You can // disable this behavior with `--watch=always`. if (args['--watch'] !== 'always') { process.stdin.on('end', () => { - Promise.all(cleanupWatchers.map((fn) => fn())).then( - () => process.exit(0), - () => process.exit(1), - ) + Promise.all(cleanupWatchers.map((fn) => fn())) + .then(() => eventBatches?.close()) + .then( + () => process.exit(0), + () => process.exit(1), + ) }) } @@ -515,6 +538,7 @@ export async function handle(args: Result>) { } await write(output, map, args, I) + finishInitialBuild() let end = process.hrtime.bigint() if (!args['--silent']) eprintln(`Done in ${formatDuration(end - start)}`) @@ -693,9 +717,29 @@ async function loadWatcher(): Promise { } } -async function createWatchers(dirs: string[], cb: (files: string[]) => void) { - let watcher = await loadWatcher() +type WatchEvent = { type: 'create' | 'update' | 'delete'; path: string } +type WatcherBackend = { + subscribe( + directory: string, + callback: (error: Error | null, events: WatchEvent[]) => Promise, + ): Promise<{ unsubscribe(): void | Promise }> +} +export async function createWatchers( + dirs: string[], + cb: (files: string[]) => Promise, + batches: SerialBatches, + watcher?: WatcherBackend, + lstat: (path: string) => Promise> = fs.lstat, +) { + if (!watcher) { + let nativeWatcher = await loadWatcher() + watcher = { + subscribe(directory, callback) { + return nativeWatcher.subscribe(directory, callback) + }, + } + } // Remove any directories that are children of an already watched directory. // If we don't we may not get notified of certain filesystem events regardless // of whether or not they are for the directory that is duplicated. @@ -730,6 +774,7 @@ async function createWatchers(dirs: string[], cb: (files: string[]) => void) { // Keep track of the debounce queue to avoid multiple rebuilds. let debounceQueue = new Disposables() + let activeCallbacks = new Set>() // A changed file can be watched by multiple watchers, but we only want to // handle the file once. We debounce the handle function with the collected @@ -740,50 +785,60 @@ async function createWatchers(dirs: string[], cb: (files: string[]) => void) { // Setup a new macrotask to handle the files in batch. debounceQueue.queueMacrotask(() => { - cb(Array.from(files)) + let batch = Array.from(files) files.clear() + void batches.push(batch) }) } // Setup a watcher for every directory. for (let dir of dirs) { - let { unsubscribe } = await watcher.subscribe(dir, async (err, events) => { - // Whenever an error occurs we want to let the user know about it but we - // want to keep watching for changes. - if (err) { - console.error(err) - return - } + let { unsubscribe } = await watcher.subscribe(dir, (err, events) => { + let callback = (async () => { + // Whenever an error occurs we want to let the user know about it but we + // want to keep watching for changes. + if (err) { + console.error(err) + return + } - await Promise.all( - events.map(async (event) => { - // When a file is deleted, a rebuild should be triggered such that we - // can figure out whether this file must trigger a fresh build or not. - // - // If it must trigger a fresh build, then we will temporarily end up - // in a broken state, but an error will be shown to the user. Once the - // user resolves the issue, the CLI will recover. - if (event.type === 'delete') { + await Promise.all( + events.map(async (event) => { + // When a file is deleted, a rebuild should be triggered such that we + // can figure out whether this file must trigger a fresh build or not. + // + // If it must trigger a fresh build, then we will temporarily end up + // in a broken state, but an error will be shown to the user. Once the + // user resolves the issue, the CLI will recover. + if (event.type === 'delete') { + files.add(event.path) + return + } + + // Ignore directory changes. We only care about file changes + let stats: Stats | null = null + try { + stats = (await lstat(event.path)) as Stats + } catch {} + if (!stats?.isFile() && !stats?.isSymbolicLink()) { + return + } + + // Track the changed file. files.add(event.path) - return - } + }), + ) - // Ignore directory changes. We only care about file changes - let stats: Stats | null = null - try { - stats = await fs.lstat(event.path) - } catch {} - if (!stats?.isFile() && !stats?.isSymbolicLink()) { - return - } + // Handle the tracked files at some point in the future. + await enqueueCallback() + })() - // Track the changed file. - files.add(event.path) - }), + activeCallbacks.add(callback) + void callback.then( + () => activeCallbacks.delete(callback), + () => activeCallbacks.delete(callback), ) - - // Handle the tracked files at some point in the future. - await enqueueCallback() + return callback }) // Ensure we cleanup the watcher when we're done. @@ -791,12 +846,29 @@ async function createWatchers(dirs: string[], cb: (files: string[]) => void) { } // Cleanup - return async () => { - await watchers.dispose() - await debounceQueue.dispose() + return { + callback: cb, + cleanup: async () => { + await watchers.dispose() + await Promise.all(activeCallbacks) + await debounceQueue.dispose() + if (files.size > 0) { + let batch = Array.from(files) + files.clear() + void batches.push(batch) + } + }, } } +export function filterChangedFiles( + files: string[], + output: string | null, + map: boolean | string, +): string[] { + return files.filter((file) => file !== output && file !== map) +} + function getRebuildStrategy( files: string[], fullRebuildPaths: string[], diff --git a/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts b/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts new file mode 100644 index 000000000..a37455bf8 --- /dev/null +++ b/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts @@ -0,0 +1,155 @@ +import { expect, it } from 'vitest' +import { serializeBatches } from './serial-batches' + +it('serializes callbacks and coalesces batches received while one is running', async () => { + let releaseFirst!: () => void + let firstCanFinish = new Promise((resolve) => (releaseFirst = resolve)) + let batches: string[][] = [] + let active = 0 + let maxActive = 0 + + let batchesQueue = serializeBatches(async (batch) => { + batches.push(batch) + active++ + maxActive = Math.max(maxActive, active) + + if (batches.length === 1) { + await firstCanFinish + } + + active-- + }) + + let first = batchesQueue.push(['a']) + await Promise.resolve() + + let second = batchesQueue.push(['b']) + let third = batchesQueue.push(['c']) + + expect(batches).toEqual([['a']]) + expect(maxActive).toBe(1) + + releaseFirst() + await Promise.all([first, second, third]) + + expect(batches).toEqual([['a'], ['b', 'c']]) + expect(maxActive).toBe(1) +}) + +it('drains accepted batches and ignores new work after close', async () => { + let release!: () => void + let canFinish = new Promise((resolve) => (release = resolve)) + let calls: string[][] = [] + let batchesQueue = serializeBatches(async (batch) => { + calls.push(batch) + await canFinish + }) + + void batchesQueue.push(['accepted']) + await Promise.resolve() + let closing = batchesQueue.close() + await batchesQueue.push(['late']) + release() + await closing + + expect(calls).toEqual([['accepted']]) +}) + +it('holds early watcher events until the initial build is complete', async () => { + let finishInitialBuild!: () => void + let initialBuild = new Promise((resolve) => (finishInitialBuild = resolve)) + let calls: string[][] = [] + let batchesQueue = serializeBatches(async (batch) => { + calls.push(batch) + }, initialBuild) + + let earlyEvent = batchesQueue.push(['changed-during-initial-build']) + await Promise.resolve() + expect(calls).toEqual([]) + + finishInitialBuild() + await earlyEvent + expect(calls).toEqual([['changed-during-initial-build']]) +}) + +it('reports callback failures and continues draining accepted batches', async () => { + let releaseFirst!: () => void + let firstCanFail = new Promise((resolve) => (releaseFirst = resolve)) + let calls: string[][] = [] + let errors: unknown[] = [] + let batchesQueue = serializeBatches( + async (batch) => { + calls.push(batch) + if (calls.length === 1) { + await firstCanFail + throw new Error('rebuild failed') + } + }, + Promise.resolve(), + (error) => errors.push(error), + ) + + let first = batchesQueue.push(['first']) + await Promise.resolve() + let second = batchesQueue.push(['accepted-during-first']) + releaseFirst() + await Promise.all([first, second]) + + expect(calls).toEqual([['first'], ['accepted-during-first']]) + expect(errors).toHaveLength(1) + expect(errors[0]).toEqual(new Error('rebuild failed')) +}) + +it('drains work queued by a callback promise reaction', async () => { + let finishFirst!: () => void + let firstCallback = new Promise((resolve) => (finishFirst = resolve)) + let calls: string[][] = [] + let queue = serializeBatches((batch) => { + calls.push(batch) + return calls.length === 1 ? firstCallback : Promise.resolve() + }) + + let first = queue.push(['first']) + let reaction = firstCallback.then(() => queue.push(['queued-by-reaction'])) + finishFirst() + await Promise.all([first, reaction]) + + expect(calls).toEqual([['first'], ['queued-by-reaction']]) +}) + +it('deduplicates repeated items across pending batches', async () => { + let release!: () => void + let blocked = new Promise((resolve) => (release = resolve)) + let calls: string[][] = [] + let queue = serializeBatches(async (batch) => { + calls.push(batch) + if (calls.length === 1) await blocked + }) + + void queue.push(['first']) + await Promise.resolve() + void queue.push(['same', 'same']) + void queue.push(['same']) + release() + await queue.close() + + expect(calls).toEqual([['first'], ['same']]) +}) + +it('reports a rejected initial barrier once and settles', async () => { + let errors: unknown[] = [] + let calls: string[][] = [] + let queue = serializeBatches( + async (batch) => { + calls.push(batch) + }, + Promise.reject(new Error('initial build failed')), + (error) => errors.push(error), + ) + + await queue.push(['change']) + await queue.close() + + expect(calls).toEqual([]) + expect(errors).toEqual([new Error('initial build failed')]) +}) diff --git a/packages/@tailwindcss-cli/src/utils/serial-batches.ts b/packages/@tailwindcss-cli/src/utils/serial-batches.ts new file mode 100644 index 000000000..ec85b69eb --- /dev/null +++ b/packages/@tailwindcss-cli/src/utils/serial-batches.ts @@ -0,0 +1,65 @@ +export interface SerialBatches { + push(batch: T[]): Promise + close(): Promise +} + +export function serializeBatches( + callback: (batch: T[]) => Promise, + startAfter: Promise = Promise.resolve(), + onError: (error: unknown) => void = console.error, +): SerialBatches { + let pending = new Set() + let inFlight: Promise | null = null + let closed = false + let startErrorReported = false + + function report(error: unknown) { + try { + onError(error) + } catch {} + } + + function startDrain(): Promise { + inFlight = (async () => { + try { + await startAfter + } catch (error) { + if (!startErrorReported) { + startErrorReported = true + report(error) + } + pending.clear() + return + } + while (pending.size > 0) { + let next = Array.from(pending) + pending.clear() + try { + await callback(next) + } catch (error) { + report(error) + } + } + })().finally(() => { + inFlight = null + if (pending.size > 0) return startDrain() + }) + + return inFlight + } + + function push(batch: T[]): Promise { + if (closed) return Promise.resolve() + for (let item of batch) pending.add(item) + + return inFlight ?? startDrain() + } + + return { + push, + async close() { + closed = true + await inFlight + }, + } +} From 2b9a5c411281ab12517f030d7641b3d94bb0e731 Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 15:24:27 +0200 Subject: [PATCH 2/7] Add a regression test for the newest change winning MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The existing tests cover the scheduling: serial-batches asserts that only one rebuild runs at a time, and the watcher tests cover event filtering and shutdown flushing. Neither asserts the outcome the bug was reported as — that what ends up written is the newest change. This drives the watcher wiring with an older change whose rebuild is slow and a newer change whose rebuild is fast, and asserts the newer result is written last. Against the previous fire-and-forget behaviour the two rebuilds overlap and the stale one lands on top, so the test fails. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/commands/build/index.test.ts | 38 +++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts index dbc0cc42d..6b4e95d43 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.test.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -76,3 +76,41 @@ it('waits for an entered watcher callback before shutdown flushes changes', asyn await queue.close() expect(calls).toEqual([['delayed-change']]) }) + +it('writes the newest change last when an earlier rebuild is slower', async () => { + // The reported bug: a rebuild for an older change finished *after* the rebuild + // for a newer one and overwrote it, leaving stale CSS on disk. The serial + // batches tests pin the scheduling on its own; this pins the outcome through + // the watcher wiring, which is the shape the bug was reported in. + let written: string[] = [] + let releaseSlowRebuild!: () => void + let slowRebuildCanFinish = new Promise((resolve) => (releaseSlowRebuild = resolve)) + + let queue = serializeBatches(async (files) => { + // Make the *first* rebuild the slow one. Without serialization the second + // rebuild finishes first and this stale result lands on top of it. + if (written.length === 0) await slowRebuildCanFinish + written.push(files.at(-1)!) + }) + let fake = fakeWatcher() + await createWatchers( + ['/watch'], + async () => {}, + queue, + fake.watcher, + async () => ({ + isFile: () => true, + isSymbolicLink: () => false, + }), + ) + + await fake.callbacks[0](null, [{ type: 'update', path: 'older-change' }]) + await nextTask() + await fake.callbacks[0](null, [{ type: 'update', path: 'newer-change' }]) + await nextTask() + + releaseSlowRebuild() + await queue.close() + + expect(written).toEqual(['older-change', 'newer-change']) +}) From d346922c115bb9bc6ef8246b32f6570af71ccae0 Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 15:50:20 +0200 Subject: [PATCH 3/7] Make only the first rebuild wait in the newest-change test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The test gated the slow rebuild on `written.length === 0`, but the first rebuild is suspended at that await, so the second rebuild saw a length of 0 too and waited on the same promise. Both then resolved in call order and the assertion held even with no serialization at all — the test could not fail. Count rebuilds instead, and wait for the first one to have actually started before submitting the second change. Verified both directions against a non-serialized queue: the old shape passed, the new one produces ['newer-change', 'older-change'] and fails. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/commands/build/index.test.ts | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts index 6b4e95d43..6cddee6e0 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.test.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -83,13 +83,20 @@ it('writes the newest change last when an earlier rebuild is slower', async () = // batches tests pin the scheduling on its own; this pins the outcome through // the watcher wiring, which is the shape the bug was reported in. let written: string[] = [] + let rebuildCount = 0 + let startFirstRebuild!: () => void + let firstRebuildStarted = new Promise((resolve) => (startFirstRebuild = resolve)) let releaseSlowRebuild!: () => void let slowRebuildCanFinish = new Promise((resolve) => (releaseSlowRebuild = resolve)) let queue = serializeBatches(async (files) => { - // Make the *first* rebuild the slow one. Without serialization the second - // rebuild finishes first and this stale result lands on top of it. - if (written.length === 0) await slowRebuildCanFinish + // Only the *first* rebuild is slow. Counting rebuilds rather than writes + // matters: the first rebuild is suspended below, so a write-count check + // would also suspend the second one and the test would pass unserialized. + if (rebuildCount++ === 0) { + startFirstRebuild() + await slowRebuildCanFinish + } written.push(files.at(-1)!) }) let fake = fakeWatcher() @@ -105,7 +112,7 @@ it('writes the newest change last when an earlier rebuild is slower', async () = ) await fake.callbacks[0](null, [{ type: 'update', path: 'older-change' }]) - await nextTask() + await firstRebuildStarted await fake.callbacks[0](null, [{ type: 'update', path: 'newer-change' }]) await nextTask() From 0bbf270efcb8096d9af9c626733d786b6318e0fb Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 16:03:22 +0200 Subject: [PATCH 4/7] Don't close the queue while a watcher swap is in flight MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A rebuild spliced the old cleanups out of `cleanupWatchers` and only pushed the new generation back after awaiting them. If stdin closed inside that window, shutdown saw an empty list, resolved immediately and closed the queue — so the files the old generation flushed on its way out were pushed into a closed queue and ignored, and the process exited 0 with stale CSS. Register the new generation before awaiting the old one, so the list is never empty, and track the in-flight swap so shutdown waits for it before closing. The ordering now lives in `shutdownWatchMode` so it can be asserted directly: the queue must still be open while a swap is flushing. Against the previous ordering that assertion fails — the queue is already closed and the flush has not landed. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/commands/build/index.test.ts | 28 ++++++++++++- .../src/commands/build/index.ts | 41 ++++++++++++++----- 2 files changed, 58 insertions(+), 11 deletions(-) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts index 6cddee6e0..b8ea3c5b6 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.test.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -1,6 +1,6 @@ import { expect, it } from 'vitest' import { serializeBatches } from '../../utils/serial-batches' -import { createWatchers, filterChangedFiles } from './index' +import { createWatchers, filterChangedFiles, shutdownWatchMode } from './index' type WatchEvent = { type: 'create' | 'update' | 'delete'; path: string } type WatchCallback = (error: Error | null, events: WatchEvent[]) => Promise @@ -121,3 +121,29 @@ it('writes the newest change last when an earlier rebuild is slower', async () = expect(written).toEqual(['older-change', 'newer-change']) }) + +it('does not close the queue while a watcher swap is still flushing', async () => { + // A rebuild swaps the watcher generation, and the old generation flushes what + // it collected as it is torn down. If shutdown closes the queue first, those + // files land in a closed queue and the process exits with stale CSS. + let closed = false + let flushed: string[] = [] + let finishSwap!: () => void + let swap = new Promise((resolve) => (finishSwap = resolve)).then(() => { + flushed.push('collected-during-swap') + }) + + let shutdown = shutdownWatchMode(swap, [], { + async close() { + closed = true + }, + }) + await nextTask() + expect(closed).toBe(false) + + finishSwap() + await shutdown + + expect(flushed).toEqual(['collected-during-swap']) + expect(closed).toBe(true) +}) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.ts b/packages/@tailwindcss-cli/src/commands/build/index.ts index 950692939..a2792b856 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.ts @@ -327,6 +327,10 @@ export async function handle(args: Result>) { let [compiler, scanner] = await handleError(() => createCompiler(input, I)) let cleanupWatchers: (() => Promise)[] = [] + // A rebuild swaps the watcher generation. Shutdown must not close the queue + // while that swap is in flight, or the files the old generation flushes are + // pushed into a closed queue and silently dropped. + let watcherSwap: Promise = Promise.resolve() let finishInitialBuild!: () => void let initialBuildFinished = new Promise((resolve) => (finishInitialBuild = resolve)) let eventBatches: SerialBatches | null = null @@ -409,12 +413,14 @@ export async function handle(args: Result>) { ) DEBUG && I.end('Setup new watchers') - // Clear old watchers + // Clear old watchers. Register the new generation *before* awaiting + // the old one, so shutdown never observes an empty cleanup list. DEBUG && I.start('Cleanup old watchers') - await Promise.all(cleanupWatchers.splice(0).map((cleanup) => cleanup())) - DEBUG && I.end('Cleanup old watchers') - + let previousCleanups = cleanupWatchers.splice(0) cleanupWatchers.push(newWatchers.cleanup) + watcherSwap = Promise.all(previousCleanups.map((cleanup) => cleanup())) + await watcherSwap + DEBUG && I.end('Cleanup old watchers') // Re-compile the CSS DEBUG && I.start('Build CSS') @@ -509,12 +515,10 @@ export async function handle(args: Result>) { // disable this behavior with `--watch=always`. if (args['--watch'] !== 'always') { process.stdin.on('end', () => { - Promise.all(cleanupWatchers.map((fn) => fn())) - .then(() => eventBatches?.close()) - .then( - () => process.exit(0), - () => process.exit(1), - ) + shutdownWatchMode(watcherSwap, cleanupWatchers, eventBatches).then( + () => process.exit(0), + () => process.exit(1), + ) }) } @@ -706,6 +710,23 @@ export async function handle(args: Result>) { // Load `@parcel/watcher` lazily so a missing or broken native binding only // affects `--watch` (without `--poll`), instead of crashing one-off builds and // polling mode as well. +/// Shut watch mode down in an order that cannot drop collected files. +/// +/// A rebuild swaps the watcher generation, and the old generation flushes what it +/// collected as it is torn down. Closing the queue before that flush lands means the +/// files are pushed into a closed queue and ignored, so the process exits successfully +/// with stale CSS. Wait for an in-flight swap first, then the current generation, and +/// only then close. +export async function shutdownWatchMode( + watcherSwap: Promise, + cleanups: (() => Promise)[], + batches: { close(): Promise } | null, +) { + await watcherSwap + await Promise.all(cleanups.map((cleanup) => cleanup())) + await batches?.close() +} + async function loadWatcher(): Promise { try { return (await import('@parcel/watcher')).default From 26fcae0e1afa2916f0fcc26e342645a3e7e00b59 Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 16:18:31 +0200 Subject: [PATCH 5/7] Re-read the watcher swap while shutting down MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Shutdown captured the swap promise once, so it waited for whichever swap was current when stdin closed. A rebuild already in flight can start its own swap after that point, and the queue was then closed while the later one was still flushing — the same dropped changes, one interleaving further out. Read the current swap each time round instead, until it stops changing. The regression test drives that interleaving directly: it replaces the swap while shutdown is already waiting, and asserts the queue stays open until the later one has flushed. Against capture-by-value the queue is already closed with only the first swap's work landed. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/commands/build/index.test.ts | 37 ++++++++++++++++++- .../src/commands/build/index.ts | 14 +++++-- 2 files changed, 47 insertions(+), 4 deletions(-) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts index b8ea3c5b6..6e02dc0f6 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.test.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -133,7 +133,7 @@ it('does not close the queue while a watcher swap is still flushing', async () = flushed.push('collected-during-swap') }) - let shutdown = shutdownWatchMode(swap, [], { + let shutdown = shutdownWatchMode(() => swap, [], { async close() { closed = true }, @@ -147,3 +147,38 @@ it('does not close the queue while a watcher swap is still flushing', async () = expect(flushed).toEqual(['collected-during-swap']) expect(closed).toBe(true) }) + +it('waits for a watcher swap that starts after shutdown begins', async () => { + // A rebuild already in flight can swap watchers while we are shutting down. + // Reading the swap once captures whichever was current when stdin closed, and + // closes the queue while the later one is still flushing. + let closed = false + let flushed: string[] = [] + let finishFirstSwap!: () => void + let firstSwap = new Promise((resolve) => (finishFirstSwap = resolve)).then(() => { + flushed.push('first-swap') + }) + let finishSecondSwap!: () => void + let secondSwap = new Promise((resolve) => (finishSecondSwap = resolve)).then(() => { + flushed.push('second-swap') + }) + + let currentSwap: Promise = firstSwap + let shutdown = shutdownWatchMode(() => currentSwap, [], { + async close() { + closed = true + }, + }) + + // The in-flight rebuild starts its own swap while shutdown is already waiting. + currentSwap = secondSwap + finishFirstSwap() + await nextTask() + expect(closed).toBe(false) + + finishSecondSwap() + await shutdown + + expect(flushed).toEqual(['first-swap', 'second-swap']) + expect(closed).toBe(true) +}) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.ts b/packages/@tailwindcss-cli/src/commands/build/index.ts index a2792b856..da7a9a218 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.ts @@ -515,7 +515,7 @@ export async function handle(args: Result>) { // disable this behavior with `--watch=always`. if (args['--watch'] !== 'always') { process.stdin.on('end', () => { - shutdownWatchMode(watcherSwap, cleanupWatchers, eventBatches).then( + shutdownWatchMode(() => watcherSwap, cleanupWatchers, eventBatches).then( () => process.exit(0), () => process.exit(1), ) @@ -718,11 +718,19 @@ export async function handle(args: Result>) { /// with stale CSS. Wait for an in-flight swap first, then the current generation, and /// only then close. export async function shutdownWatchMode( - watcherSwap: Promise, + pendingSwap: () => Promise, cleanups: (() => Promise)[], batches: { close(): Promise } | null, ) { - await watcherSwap + // A rebuild already in flight can start its own swap while we are shutting + // down, so read the current one each time round rather than capturing it + // once. Capturing it once waits for the swap that happened to be current when + // stdin closed and closes the queue while a later one is still flushing. + let awaited: Promise | undefined + while (awaited !== pendingSwap()) { + awaited = pendingSwap() + await awaited + } await Promise.all(cleanups.map((cleanup) => cleanup())) await batches?.close() } From 5302f308fca8d64c7813ee0c8d3c0afd12d14469 Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 16:26:20 +0200 Subject: [PATCH 6/7] Drain the queue on shutdown instead of tracking watcher swaps MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Tracking the swap promise was the wrong shape. It could always be read at a moment when the in-flight rebuild had not assigned its swap yet, so shutdown saw the old one settle, closed the queue, and dropped whatever the later generation flushed. Each fix narrowed that window without closing it. A full rebuild runs inside the queue, so the queue already knows when one is in flight — including a rebuild that has not yet replaced its watchers. Add `drain()`, wait for it before tearing the watchers down, and close only after their flush has landed. That removes the swap bookkeeping entirely. Both orderings are pinned: cleanup must not run while a rebuild is in flight, and what the watchers flush on the way out must still be processed. Against no-drain the cleanup runs early; against closing first the flush is ignored. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/commands/build/index.test.ts | 88 +++++++++---------- .../src/commands/build/index.ts | 32 ++----- .../src/utils/serial-batches.ts | 11 +++ 3 files changed, 60 insertions(+), 71 deletions(-) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.test.ts b/packages/@tailwindcss-cli/src/commands/build/index.test.ts index 6e02dc0f6..aeca6c9ec 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.test.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.test.ts @@ -122,63 +122,55 @@ it('writes the newest change last when an earlier rebuild is slower', async () = expect(written).toEqual(['older-change', 'newer-change']) }) -it('does not close the queue while a watcher swap is still flushing', async () => { - // A rebuild swaps the watcher generation, and the old generation flushes what - // it collected as it is torn down. If shutdown closes the queue first, those - // files land in a closed queue and the process exits with stale CSS. - let closed = false - let flushed: string[] = [] - let finishSwap!: () => void - let swap = new Promise((resolve) => (finishSwap = resolve)).then(() => { - flushed.push('collected-during-swap') +it('waits for an in-flight rebuild before tearing down the watchers', async () => { + // A full rebuild runs inside the queue and swaps the watcher generation while + // it does. Tearing the watchers down first races that swap. + let order: string[] = [] + let finishRebuild!: () => void + let rebuildDone = new Promise((resolve) => (finishRebuild = resolve)) + let queue = serializeBatches(async () => { + order.push('rebuild:start') + await rebuildDone + order.push('rebuild:end') }) - let shutdown = shutdownWatchMode(() => swap, [], { - async close() { - closed = true - }, - }) + void queue.push(['change']) await nextTask() - expect(closed).toBe(false) - finishSwap() + let shutdown = shutdownWatchMode( + [ + async () => { + order.push('cleanup') + }, + ], + queue, + ) + await nextTask() + expect(order).toEqual(['rebuild:start']) + + finishRebuild() await shutdown - expect(flushed).toEqual(['collected-during-swap']) - expect(closed).toBe(true) + expect(order).toEqual(['rebuild:start', 'rebuild:end', 'cleanup']) }) -it('waits for a watcher swap that starts after shutdown begins', async () => { - // A rebuild already in flight can swap watchers while we are shutting down. - // Reading the swap once captures whichever was current when stdin closed, and - // closes the queue while the later one is still flushing. - let closed = false - let flushed: string[] = [] - let finishFirstSwap!: () => void - let firstSwap = new Promise((resolve) => (finishFirstSwap = resolve)).then(() => { - flushed.push('first-swap') - }) - let finishSecondSwap!: () => void - let secondSwap = new Promise((resolve) => (finishSecondSwap = resolve)).then(() => { - flushed.push('second-swap') +it('processes what the watchers flush on the way out', async () => { + // The watchers flush what they collected as they are torn down. That has to + // land in a queue that is still open, or it is dropped and we exit as if all + // was well. + let processed: string[][] = [] + let queue = serializeBatches(async (files) => { + processed.push(files) }) - let currentSwap: Promise = firstSwap - let shutdown = shutdownWatchMode(() => currentSwap, [], { - async close() { - closed = true - }, - }) + await shutdownWatchMode( + [ + async () => { + void queue.push(['flushed-on-shutdown']) + }, + ], + queue, + ) - // The in-flight rebuild starts its own swap while shutdown is already waiting. - currentSwap = secondSwap - finishFirstSwap() - await nextTask() - expect(closed).toBe(false) - - finishSecondSwap() - await shutdown - - expect(flushed).toEqual(['first-swap', 'second-swap']) - expect(closed).toBe(true) + expect(processed).toEqual([['flushed-on-shutdown']]) }) diff --git a/packages/@tailwindcss-cli/src/commands/build/index.ts b/packages/@tailwindcss-cli/src/commands/build/index.ts index da7a9a218..9b82b4fe7 100644 --- a/packages/@tailwindcss-cli/src/commands/build/index.ts +++ b/packages/@tailwindcss-cli/src/commands/build/index.ts @@ -327,10 +327,6 @@ export async function handle(args: Result>) { let [compiler, scanner] = await handleError(() => createCompiler(input, I)) let cleanupWatchers: (() => Promise)[] = [] - // A rebuild swaps the watcher generation. Shutdown must not close the queue - // while that swap is in flight, or the files the old generation flushes are - // pushed into a closed queue and silently dropped. - let watcherSwap: Promise = Promise.resolve() let finishInitialBuild!: () => void let initialBuildFinished = new Promise((resolve) => (finishInitialBuild = resolve)) let eventBatches: SerialBatches | null = null @@ -418,8 +414,7 @@ export async function handle(args: Result>) { DEBUG && I.start('Cleanup old watchers') let previousCleanups = cleanupWatchers.splice(0) cleanupWatchers.push(newWatchers.cleanup) - watcherSwap = Promise.all(previousCleanups.map((cleanup) => cleanup())) - await watcherSwap + await Promise.all(previousCleanups.map((cleanup) => cleanup())) DEBUG && I.end('Cleanup old watchers') // Re-compile the CSS @@ -515,7 +510,7 @@ export async function handle(args: Result>) { // disable this behavior with `--watch=always`. if (args['--watch'] !== 'always') { process.stdin.on('end', () => { - shutdownWatchMode(() => watcherSwap, cleanupWatchers, eventBatches).then( + shutdownWatchMode(cleanupWatchers, eventBatches).then( () => process.exit(0), () => process.exit(1), ) @@ -712,25 +707,16 @@ export async function handle(args: Result>) { // polling mode as well. /// Shut watch mode down in an order that cannot drop collected files. /// -/// A rebuild swaps the watcher generation, and the old generation flushes what it -/// collected as it is torn down. Closing the queue before that flush lands means the -/// files are pushed into a closed queue and ignored, so the process exits successfully -/// with stale CSS. Wait for an in-flight swap first, then the current generation, and -/// only then close. +/// A full rebuild runs inside the queue and swaps the watcher generation while it +/// does, so draining the queue is what waits for a rebuild to finish — including +/// one that has not yet replaced its watchers. Only then is it safe to tear the +/// watchers down, because the files they flush on the way out have to land in a +/// queue that is still open. Closing first drops them and exits with stale CSS. export async function shutdownWatchMode( - pendingSwap: () => Promise, cleanups: (() => Promise)[], - batches: { close(): Promise } | null, + batches: { drain(): Promise; close(): Promise } | null, ) { - // A rebuild already in flight can start its own swap while we are shutting - // down, so read the current one each time round rather than capturing it - // once. Capturing it once waits for the swap that happened to be current when - // stdin closed and closes the queue while a later one is still flushing. - let awaited: Promise | undefined - while (awaited !== pendingSwap()) { - awaited = pendingSwap() - await awaited - } + await batches?.drain() await Promise.all(cleanups.map((cleanup) => cleanup())) await batches?.close() } diff --git a/packages/@tailwindcss-cli/src/utils/serial-batches.ts b/packages/@tailwindcss-cli/src/utils/serial-batches.ts index ec85b69eb..5ee41db7e 100644 --- a/packages/@tailwindcss-cli/src/utils/serial-batches.ts +++ b/packages/@tailwindcss-cli/src/utils/serial-batches.ts @@ -1,5 +1,9 @@ export interface SerialBatches { push(batch: T[]): Promise + /// Wait for in-flight work to finish, without closing. Callers that need the + /// queue to still accept work afterwards — a shutdown that has yet to flush + /// what the watchers collected — drain first and close last. + drain(): Promise close(): Promise } @@ -57,6 +61,13 @@ export function serializeBatches( return { push, + async drain() { + // `inFlight` is re-chained by finalization when work arrived mid-drain, so + // loop until it settles rather than awaiting whichever promise is current. + while (inFlight) { + await inFlight + } + }, async close() { closed = true await inFlight From 3148fad1a1cc9bbdd02c90911c98b4ee6746e24b Mon Sep 17 00:00:00 2001 From: Michael Glass Date: Mon, 31 Aug 2026 16:38:19 +0200 Subject: [PATCH 7/7] Test drain directly drain() is public on SerialBatches but was only exercised through shutdown. Cover it on its own: it waits for in-flight work and leaves the queue open, so a shutdown can drain before the watchers flush and still have that work accepted, and it returns immediately when nothing is running. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Lcj4iQ3fBxMwAu2rf4zLbC --- .../src/utils/serial-batches.test.ts | 33 +++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts b/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts index a37455bf8..c7adb77b9 100644 --- a/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts +++ b/packages/@tailwindcss-cli/src/utils/serial-batches.test.ts @@ -153,3 +153,36 @@ it('reports a rejected initial barrier once and settles', async () => { expect(calls).toEqual([]) expect(errors).toEqual([new Error('initial build failed')]) }) + +it('drains in-flight work but stays open for more', async () => { + let release!: () => void + let canFinish = new Promise((resolve) => (release = resolve)) + let calls: string[][] = [] + let queue = serializeBatches(async (batch) => { + calls.push(batch) + if (calls.length === 1) await canFinish + }) + + void queue.push(['first']) + await Promise.resolve() + + let drained = queue.drain() + release() + await drained + expect(calls).toEqual([['first']]) + + // Draining must not close the queue — a shutdown drains before the watchers + // have flushed what they collected, and that work still has to be accepted. + await queue.push(['after-drain']) + expect(calls).toEqual([['first'], ['after-drain']]) +}) + +it('drain returns immediately when nothing is in flight', async () => { + let calls: string[][] = [] + let queue = serializeBatches(async (batch) => { + calls.push(batch) + }) + + await queue.drain() + expect(calls).toEqual([]) +})