Compare commits

..
Author SHA1 Message Date
sudacode 354e3ea2e5 test(dictionary): pin concurrent snapshot writers to one surviving snapshot
Give each concurrent writer a distinct title, length, and term text so the
assertion identifies which writer's snapshot survived instead of only
checking that the file parses with the expected entry count.
2026-08-17 11:10:38 -07:00
sudacode 9df2f37d4b fix(dictionary): isolate concurrent snapshot writes and shorten zip build blocks
The streamed snapshot writer keyed its temp file on the pid alone, so two
overlapping writes for the same media (a manual generate racing auto-sync)
streamed into one file and tore it; add a per-write sequence suffix.

Term banks were stringified 10k entries at a time, measured at ~38MB and
~135ms per bank on a real merged dictionary. Halving that block matters
more than the write itself: at 2k entries the longest event-loop stall in
a merged build drops from ~135ms to 28ms.
2026-08-17 02:46:27 -07:00
sudacode 1228ebe622 fix(dictionary): stop character dictionary IO from blocking the main process
Generating and importing a large character dictionary froze the whole
app long enough for the compositor to raise its application-not-
responding dialog over the player. Multi-hundred-MB snapshot JSONs and
the merged archive were read, written, and zipped synchronously on the
main process, and the character image / name-candidate caches re-read
every cached snapshot synchronously inside a lookup whenever the
snapshot directory changed.

- snapshot reads/writes are async; writes stream in slices and rename
  into place so a crash or concurrent writer cannot tear a snapshot
- buildDictionaryZip yields between ~8MB slices and CRC32 uses the
  native zlib implementation
- the image and name-candidate lookup caches rebuild in the background
  and serve the previous index while the rebuild runs

Worst main-thread stall over a 1.4GB snapshot set drops from 8s+ to
under 700ms.
2026-08-17 01:45:54 -07:00
sudacode ec7f0345b0 fix(notifications): spawn notify-send without AppImage library overrides
Electron AppImages export LD_LIBRARY_PATH pointing at bundled libraries
whose stale libnotify kills the system notify-send with a symbol lookup
error, permanently disabling in-place replacement and forcing the
flickering Electron close-and-reopen fallback. Drop the override from
the child environment so the system binary resolves its own libraries.
2026-08-17 01:45:47 -07:00
sudacode 00b1b79bf4 fix(mpv): recover from stalled IPC connects (#204) 2026-08-16 22:58:28 -07:00
23 changed files with 872 additions and 157 deletions
@@ -0,0 +1,5 @@
type: fixed
area: dictionary
- Character dictionary generation, merged rebuilds, and imports no longer freeze the app (and trigger the compositor's "application not responding" dialog) on large dictionaries; snapshot reads/writes, archive building, and the character image/name lookup caches now do their heavy work off the UI's critical path.
- Desktop progress notifications now update in place on Linux AppImage installs too: the AppImage's bundled libraries broke the system notify-send helper, which silently forced the flickering close-and-reopen notification fallback.
+4
View File
@@ -0,0 +1,4 @@
type: fixed
area: overlay
- Fixed the overlay getting stuck on "Overlay loading" forever when startup stalls: mpv IPC connection attempts now time out and retry, switching sockets aborts obsolete attempts, and the plugin replaces its spinner with an actionable error if overlay content is still not ready after 30 seconds.
+26
View File
@@ -7,6 +7,8 @@ local OVERLAY_RESTART_PING_MAX_ATTEMPTS = 20
local OVERLAY_LOADING_OSD_PREFIX = "Overlay loading "
local OVERLAY_LOADING_OSD_FRAMES = { "|", "/", "-", "\\" }
local OVERLAY_LOADING_OSD_REFRESH_SECONDS = 0.18
local OVERLAY_LOADING_OSD_DEADLINE_SECONDS = 30
local OVERLAY_LOADING_OSD_TIMEOUT_MESSAGE = "Overlay did not become ready; check SubMiner logs"
local AUTO_PLAY_READY_LOADING_OSD = "Loading subtitle tokenization..."
local AUTO_PLAY_READY_READY_OSD = "Subtitle tokenization ready"
local DEFAULT_AUTO_PLAY_READY_TIMEOUT_SECONDS = 30
@@ -265,10 +267,19 @@ function M.create(ctx)
state.overlay_loading_osd_timer = nil
end
local function clear_overlay_loading_osd_deadline()
local timeout = state.overlay_loading_osd_deadline
if timeout and timeout.kill then
timeout:kill()
end
state.overlay_loading_osd_deadline = nil
end
local function stop_overlay_loading_osd()
state.overlay_loading_osd_active = false
state.overlay_loading_osd_frame = 1
clear_overlay_loading_osd_timer()
clear_overlay_loading_osd_deadline()
end
local function start_overlay_loading_osd()
@@ -291,6 +302,21 @@ function M.create(ctx)
end
end)
end
if type(mp.add_timeout) == "function" then
state.overlay_loading_osd_deadline = mp.add_timeout(OVERLAY_LOADING_OSD_DEADLINE_SECONDS, function()
if not state.overlay_loading_osd_active then
return
end
state.overlay_loading_osd_deadline = nil
stop_overlay_loading_osd()
subminer_log(
"warn",
"process",
"Overlay loading deadline expired before the app reported content ready"
)
show_osd(OVERLAY_LOADING_OSD_TIMEOUT_MESSAGE, { force = true })
end)
end
end
local function disarm_auto_play_ready_gate(options)
+1
View File
@@ -26,6 +26,7 @@ function M.new()
auto_play_ready_initial_pause_ownership_consumed = false,
overlay_loading_osd_active = false,
overlay_loading_osd_timer = nil,
overlay_loading_osd_deadline = nil,
overlay_loading_osd_frame = 1,
pending_visible_overlay_hide_timer = nil,
pending_visible_overlay_hide_generation = 0,
+53 -1
View File
@@ -130,7 +130,9 @@ local function run_plugin_scenario(config)
function mp.add_timeout(seconds, callback)
recorded.timeouts[#recorded.timeouts + 1] = seconds
local delay = tonumber(seconds) or 0
local timeout = {
seconds = delay,
killed = false,
callback = callback,
}
@@ -138,7 +140,6 @@ local function run_plugin_scenario(config)
self.killed = true
end
local delay = tonumber(seconds) or 0
if callback and delay < 5 and not config.defer_timeouts then
callback()
end
@@ -514,6 +515,15 @@ local function has_timeout(timeouts, target)
return false
end
local function find_timeout_handle(recorded, target)
for _, timeout in ipairs(recorded.timeout_handles) do
if math.abs(timeout.seconds - target) < 0.0001 then
return timeout
end
end
return nil
end
local function env_has(call, target)
local env = (call and call.env) or {}
for _, value in ipairs(env) do
@@ -1636,6 +1646,8 @@ do
#recorded.periodic_timers == 1,
"auto-start visible overlay should refresh the early overlay loading OSD"
)
local overlay_loading_deadline = find_timeout_handle(recorded, 30)
assert_true(overlay_loading_deadline ~= nil, "overlay loading OSD should have a bounded deadline")
local overlay_loading_timer = recorded.periodic_timers[1]
recorded.periodic_timers[1].callback()
assert_true(
@@ -1670,6 +1682,46 @@ do
recorded.periodic_timers[1].killed == true,
"overlay loading ready should stop the early overlay loading OSD refresher"
)
assert_true(
overlay_loading_deadline.killed == true,
"overlay loading ready should cancel the bounded loading deadline"
)
end
do
local recorded, err = run_plugin_scenario({
defer_timeouts = true,
process_list = "",
option_overrides = {
binary_path = binary_path,
auto_start = "yes",
auto_start_visible_overlay = "yes",
osd_messages = false,
socket_path = "/tmp/subminer-socket",
},
input_ipc_server = "/tmp/subminer-socket",
media_title = "Random Movie",
files = {
[binary_path] = true,
},
})
assert_true(recorded ~= nil, "plugin failed to load for overlay loading deadline scenario: " .. tostring(err))
fire_event(recorded, "start-file")
local overlay_loading_deadline = find_timeout_handle(recorded, 30)
assert_true(overlay_loading_deadline ~= nil, "overlay loading deadline should be scheduled")
overlay_loading_deadline.callback()
assert_true(
recorded.periodic_timers[1].killed == true,
"overlay loading deadline should stop the loading spinner"
)
assert_true(
has_osd_message(recorded.osd, "SubMiner: Overlay did not become ready; check SubMiner logs"),
"overlay loading deadline should replace the spinner with actionable feedback"
)
assert_true(
has_log_containing(recorded.logs, "Overlay loading deadline expired"),
"overlay loading deadline should leave a diagnostic log entry"
)
end
do
+78 -1
View File
@@ -38,7 +38,15 @@ class ManualCloseSocket extends FakeSocket {
}
}
const wait = () => new Promise((resolve) => setTimeout(resolve, 0));
class HangingSocket extends FakeSocket {
override connect(path: string): void {
this.connectedPaths.push(path);
// Never emits 'connect', 'error', or 'close' on its own: models a named
// pipe dial that stalls indefinitely.
}
}
const wait = (ms = 0) => new Promise((resolve) => setTimeout(resolve, ms));
test('getMpvReconnectDelay follows existing reconnect ramp', () => {
assert.equal(getMpvReconnectDelay(0, true), 1000);
@@ -232,6 +240,75 @@ test('MpvSocketTransport.shutdown clears socket and lifecycle flags', async () =
assert.deepEqual(events, []);
});
test('MpvSocketTransport aborts a hung connect after the timeout and allows a fresh dial', async () => {
const events: string[] = [];
const errors: Error[] = [];
const sockets: HangingSocket[] = [];
const transport = new MpvSocketTransport({
socketPath: '/tmp/mpv.sock',
connectTimeoutMs: 5,
onConnect: () => {
events.push('connect');
},
onData: () => {},
onError: (error) => {
events.push('error');
errors.push(error);
},
onClose: () => {
events.push('close');
},
socketFactory: () => {
const socket = new HangingSocket();
sockets.push(socket);
return socket as unknown as net.Socket;
},
});
transport.connect();
assert.equal(transport.isConnecting, true);
await wait(20);
assert.deepEqual(events, ['error', 'close']);
assert.match(errors[0]!.message, /connect timed out/);
assert.equal(sockets[0]!.destroyed, true);
assert.equal(transport.isConnecting, false);
assert.equal(transport.isConnected, false);
transport.connect();
assert.equal(transport.isConnecting, true);
assert.equal(sockets.length, 2);
assert.equal(sockets[1]!.connectedPaths.at(0), '/tmp/mpv.sock');
transport.shutdown();
});
test('MpvSocketTransport does not fire the connect timeout after a successful connect', async () => {
const events: string[] = [];
const transport = new MpvSocketTransport({
socketPath: '/tmp/mpv.sock',
connectTimeoutMs: 5,
onConnect: () => {
events.push('connect');
},
onData: () => {},
onError: () => {
events.push('error');
},
onClose: () => {
events.push('close');
},
socketFactory: () => new FakeSocket() as unknown as net.Socket,
});
transport.connect();
await wait(20);
assert.deepEqual(events, ['connect']);
assert.equal(transport.isConnected, true);
});
test('MpvSocketTransport ignores stale socket events after shutdown and reconnect', async () => {
const events: string[] = [];
const sockets: ManualCloseSocket[] = [];
+36
View File
@@ -62,6 +62,8 @@ interface MpvSocketTransportEvents {
onClose: () => void;
}
export const MPV_CONNECT_TIMEOUT_MS = 5000;
export interface MpvSocketTransportOptions {
socketPath: string;
onConnect: () => void;
@@ -69,13 +71,16 @@ export interface MpvSocketTransportOptions {
onError: (error: Error) => void;
onClose: () => void;
socketFactory?: () => net.Socket;
connectTimeoutMs?: number;
}
export class MpvSocketTransport {
private socketPath: string;
private readonly callbacks: MpvSocketTransportEvents;
private readonly socketFactory: () => net.Socket;
private readonly connectTimeoutMs: number;
private socketRef: net.Socket | null = null;
private connectTimer: ReturnType<typeof setTimeout> | null = null;
public socket: net.Socket | null = null;
public connected = false;
public connecting = false;
@@ -83,6 +88,7 @@ export class MpvSocketTransport {
constructor(options: MpvSocketTransportOptions) {
this.socketPath = options.socketPath;
this.socketFactory = options.socketFactory ?? (() => new net.Socket());
this.connectTimeoutMs = options.connectTimeoutMs ?? MPV_CONNECT_TIMEOUT_MS;
this.callbacks = {
onConnect: options.onConnect,
onData: options.onData,
@@ -91,6 +97,31 @@ export class MpvSocketTransport {
};
}
private clearConnectTimeout(): void {
if (this.connectTimer) {
clearTimeout(this.connectTimer);
this.connectTimer = null;
}
}
// A named-pipe/socket dial that neither connects nor errors would otherwise
// latch `connecting` forever and silently block every future connect().
private armConnectTimeout(socket: net.Socket): void {
this.clearConnectTimeout();
this.connectTimer = setTimeout(() => {
this.connectTimer = null;
if (this.socketRef !== socket || this.connected) return;
this.connecting = false;
this.callbacks.onError(
new Error(`MPV IPC connect timed out after ${this.connectTimeoutMs}ms: ${this.socketPath}`),
);
// Destroying the socket emits 'close', which drives the normal
// disconnect path (including reconnect scheduling) upstream.
socket.destroy();
}, this.connectTimeoutMs);
this.connectTimer.unref?.();
}
setSocketPath(socketPath: string): void {
this.socketPath = socketPath;
}
@@ -111,6 +142,7 @@ export class MpvSocketTransport {
socket.on('connect', () => {
if (this.socketRef !== socket) return;
this.clearConnectTimeout();
this.connected = true;
this.connecting = false;
this.callbacks.onConnect();
@@ -123,6 +155,7 @@ export class MpvSocketTransport {
socket.on('error', (error: Error) => {
if (this.socketRef !== socket) return;
this.clearConnectTimeout();
this.connected = false;
this.connecting = false;
this.callbacks.onError(error);
@@ -130,12 +163,14 @@ export class MpvSocketTransport {
socket.on('close', () => {
if (this.socketRef !== socket) return;
this.clearConnectTimeout();
this.connected = false;
this.connecting = false;
this.callbacks.onClose();
});
socket.connect(this.socketPath);
this.armConnectTimeout(socket);
}
send(payload: MpvSocketMessagePayload): boolean {
@@ -149,6 +184,7 @@ export class MpvSocketTransport {
}
shutdown(): void {
this.clearConnectTimeout();
const socket = this.socketRef;
this.socketRef = null;
this.socket = null;
+127
View File
@@ -1,5 +1,6 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { EventEmitter } from 'node:events';
import {
MpvIpcClient,
MpvIpcClientDeps,
@@ -23,6 +24,18 @@ function makeDeps(overrides: Partial<MpvIpcClientProtocolDeps> = {}): MpvIpcClie
};
}
const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
async function waitFor(predicate: () => boolean, timeoutMs = 2000): Promise<void> {
const deadline = Date.now() + timeoutMs;
while (!predicate()) {
if (Date.now() >= deadline) {
throw new Error('Timed out waiting for MPV retry connection');
}
await wait(10);
}
}
function captureWarnLogs(run: () => void): string[] {
const originalWarn = console.warn;
const originalLogLevel = process.env.SUBMINER_LOG_LEVEL;
@@ -756,3 +769,117 @@ test('MpvIpcClient playNextSubtitle still auto-pauses at end while already playi
assert.equal((client as any).pendingPauseAtSubEnd, true);
assert.deepEqual(commands, [{ command: ['sub-seek', 1] }]);
});
class HangingTestSocket extends EventEmitter {
public connectedPaths: string[] = [];
public destroyed = false;
connect(path: string): void {
this.connectedPaths.push(path);
// Never resolves: models a stalled named-pipe dial.
}
write(): boolean {
return true;
}
destroy(): void {
this.destroyed = true;
}
}
class RetryTestSocket extends EventEmitter {
public connectedPaths: string[] = [];
public destroyed = false;
constructor(private readonly shouldConnect: boolean) {
super();
}
connect(path: string): void {
this.connectedPaths.push(path);
if (this.shouldConnect) {
setTimeout(() => this.emit('connect'), 0);
}
}
write(): boolean {
return true;
}
destroy(): void {
if (this.destroyed) return;
this.destroyed = true;
this.emit('close');
}
}
test('MpvIpcClient automatically retries the same socket path after a connect timeout', async () => {
const sockets: RetryTestSocket[] = [];
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
const originalLogLevel = process.env.SUBMINER_LOG_LEVEL;
const client = new MpvIpcClient(
'/tmp/mpv.sock',
makeDeps({
connectTimeoutMs: 5,
getReconnectTimer: () => reconnectTimer,
setReconnectTimer: (timer) => {
reconnectTimer = timer;
},
socketFactory: () => {
const socket = new RetryTestSocket(sockets.length > 0);
sockets.push(socket);
return socket as unknown as import('node:net').Socket;
},
}),
);
process.env.SUBMINER_LOG_LEVEL = 'error';
try {
client.connect();
await waitFor(() => client.connected);
assert.equal(sockets.length, 2);
assert.equal(sockets[0]!.destroyed, true);
assert.equal(sockets[0]!.connectedPaths.at(0), '/tmp/mpv.sock');
assert.equal(sockets[1]!.connectedPaths.at(0), '/tmp/mpv.sock');
assert.equal(client.connected, true);
} finally {
if (originalLogLevel === undefined) {
delete process.env.SUBMINER_LOG_LEVEL;
} else {
process.env.SUBMINER_LOG_LEVEL = originalLogLevel;
}
if (reconnectTimer) clearTimeout(reconnectTimer);
(client as any).transport.shutdown();
}
});
test('MpvIpcClient.setSocketPath aborts an in-flight connect so the next dial targets the new path', () => {
const sockets: HangingTestSocket[] = [];
const client = new MpvIpcClient(
'/tmp/mpv-old.sock',
makeDeps({
socketFactory: () => {
const socket = new HangingTestSocket();
sockets.push(socket);
return socket as unknown as import('node:net').Socket;
},
}),
);
client.connect();
assert.equal(sockets.length, 1);
assert.equal(sockets[0]!.connectedPaths.at(0), '/tmp/mpv-old.sock');
assert.equal((client as any).connecting, true);
client.setSocketPath('/tmp/mpv-new.sock');
assert.equal((client as any).connecting, false);
assert.equal(sockets[0]!.destroyed, true);
client.connect();
assert.equal(sockets.length, 2);
assert.equal(sockets[1]!.connectedPaths.at(0), '/tmp/mpv-new.sock');
(client as any).transport.shutdown();
});
+17 -1
View File
@@ -9,7 +9,11 @@ import {
splitMpvMessagesFromBuffer,
} from './mpv-protocol';
import { requestMpvInitialState, subscribeToMpvProperties } from './mpv-properties';
import { scheduleMpvReconnect, MpvSocketTransport } from './mpv-transport';
import {
scheduleMpvReconnect,
MpvSocketTransport,
MpvSocketTransportOptions,
} from './mpv-transport';
import { createLogger } from '../../logger';
const logger = createLogger('main:mpv');
@@ -110,6 +114,8 @@ export interface MpvIpcClientProtocolDeps {
shouldAutoLoadSecondarySubTrack?: (path: string) => boolean;
shouldQuitOnMpvShutdown?: () => boolean;
requestAppQuit?: () => void;
socketFactory?: MpvSocketTransportOptions['socketFactory'];
connectTimeoutMs?: number;
}
export interface MpvIpcClientDeps extends MpvIpcClientProtocolDeps {}
@@ -188,6 +194,8 @@ export class MpvIpcClient implements MpvClient {
this.transport = new MpvSocketTransport({
socketPath,
socketFactory: deps.socketFactory,
connectTimeoutMs: deps.connectTimeoutMs,
onConnect: () => {
this.connected = true;
this.connecting = false;
@@ -289,6 +297,14 @@ export class MpvIpcClient implements MpvClient {
previousSocketPath: this.socketPath,
socketPath,
});
if (this.connecting && !this.connected) {
// Abort the in-flight dial to the old path; otherwise the connecting
// latch turns every later connect() into a no-op while we hang on a
// stale socket.
logger.debug('Aborting in-flight MPV IPC connect for socket path change.');
this.transport.shutdown();
this.connecting = false;
}
}
this.socketPath = socketPath;
this.transport.setSocketPath(socketPath);
+17 -1
View File
@@ -1,6 +1,22 @@
import assert from 'node:assert/strict';
import test from 'node:test';
import { createNotifySendReplacer, resolveDefaultNotificationIconPath } from './notification';
import {
buildNotifySendEnv,
createNotifySendReplacer,
resolveDefaultNotificationIconPath,
} from './notification';
test('notify-send child environment drops the AppImage library-path override', () => {
const env = buildNotifySendEnv({
LD_LIBRARY_PATH: '/tmp/.mount_SubMinXXXXXX/usr/lib',
DBUS_SESSION_BUS_ADDRESS: 'unix:path=/run/user/1000/bus',
HOME: '/home/user',
});
assert.equal(env.LD_LIBRARY_PATH, undefined);
assert.equal(env.DBUS_SESSION_BUS_ADDRESS, 'unix:path=/run/user/1000/bus');
assert.equal(env.HOME, '/home/user');
});
test('default notification icon resolves packaged SubMiner asset when no per-notification icon is provided', () => {
const path = resolveDefaultNotificationIconPath({
+12 -1
View File
@@ -203,8 +203,19 @@ export function createNotifySendReplacer(
};
}
/**
* Electron AppImages export `LD_LIBRARY_PATH=<mount>/usr/lib`, whose bundled libnotify predates the
* symbols the system notify-send links against, so an inherited environment kills the child with a
* symbol lookup error before it can send anything. A system binary resolves its own libraries fine,
* so the override is dropped entirely rather than filtered.
*/
export function buildNotifySendEnv(env: NodeJS.ProcessEnv = process.env): NodeJS.ProcessEnv {
const { LD_LIBRARY_PATH: _dropped, ...rest } = env;
return rest;
}
const showLinuxReplaceableNotification = createNotifySendReplacer((args, callback) =>
execFile('notify-send', args, { timeout: 5_000 }, (error, stdout) =>
execFile('notify-send', args, { timeout: 5_000, env: buildNotifySendEnv() }, (error, stdout) =>
callback(error, stdout ?? ''),
),
);
+19 -14
View File
@@ -244,13 +244,13 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
};
};
const findCachedSnapshotForSeriesKey = (
const findCachedSnapshotForSeriesKey = async (
seriesKey: string,
fallbackSeriesKey?: string,
): CharacterDictionarySnapshot | null => {
): Promise<CharacterDictionarySnapshot | null> => {
const acceptedKeys = new Set([seriesKey, fallbackSeriesKey].filter(Boolean));
return (
readCachedSnapshots(outputDir).find((snapshot) => {
(await readCachedSnapshots(outputDir)).find((snapshot) => {
const snapshotSeriesKey = buildCharacterDictionarySeriesKey({
mediaPath: null,
mediaTitle: snapshot.mediaTitle,
@@ -293,7 +293,9 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
const cachedResolution = readCachedMediaResolution(outputDir, seriesKey);
if (cachedResolution) {
const cachedSnapshot = readSnapshot(getSnapshotPath(outputDir, cachedResolution.mediaId));
const cachedSnapshot = await readSnapshot(
getSnapshotPath(outputDir, cachedResolution.mediaId),
);
if (cachedSnapshot) {
deps.logInfo?.(
`[dictionary] cached AniList match: ${cachedSnapshot.mediaTitle} -> AniList ${cachedSnapshot.mediaId}`,
@@ -305,7 +307,7 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
}
}
const cachedSnapshot = findCachedSnapshotForSeriesKey(seriesKey, unscopedSeriesKey);
const cachedSnapshot = await findCachedSnapshotForSeriesKey(seriesKey, unscopedSeriesKey);
if (cachedSnapshot) {
writeCachedMediaResolution(outputDir, {
seriesKey,
@@ -348,7 +350,7 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
progress?: CharacterDictionarySnapshotProgressCallbacks,
): Promise<CharacterDictionarySnapshotResult> => {
const snapshotPath = getSnapshotPath(outputDir, mediaId);
const cachedSnapshot = readSnapshot(snapshotPath);
const cachedSnapshot = await readSnapshot(snapshotPath);
const refreshReason = cachedSnapshot ? getCachedSnapshotRefreshReason(cachedSnapshot) : null;
if (cachedSnapshot && refreshReason === null) {
deps.logInfo?.(`[dictionary] snapshot hit for AniList ${mediaId}`);
@@ -485,7 +487,7 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
resolvedNameSplits,
nameSplitSource,
);
writeSnapshot(snapshotPath, snapshot);
await writeSnapshot(snapshotPath, snapshot);
deps.logInfo?.(
`[dictionary] stored snapshot for AniList ${mediaId}: ${snapshot.entryCount} terms`,
);
@@ -526,19 +528,22 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
const snapshotResults = await Promise.all(
normalizedMediaIds.map((mediaId) => getOrCreateSnapshot(mediaId)),
);
const snapshots = snapshotResults.map(({ mediaId }) => {
const snapshot = readSnapshot(getSnapshotPath(outputDir, mediaId));
// Sequential on purpose: each snapshot parse is a chunk of main-thread work, so reading them
// one at a time keeps the event loop breathing between files.
const snapshots: CharacterDictionarySnapshot[] = [];
for (const { mediaId } of snapshotResults) {
const snapshot = await readSnapshot(getSnapshotPath(outputDir, mediaId));
if (!snapshot) {
throw new Error(`Missing character dictionary snapshot for AniList ${mediaId}.`);
}
return snapshot;
});
snapshots.push(snapshot);
}
const revision = buildMergedRevision(normalizedMediaIds, snapshots);
const description =
snapshots.length === 1
? `Character names from ${snapshots[0]!.mediaTitle}`
: `Character names from ${snapshots.length} recent anime`;
const { zipPath, entryCount } = buildDictionaryZip(
const { zipPath, entryCount } = await buildDictionaryZip(
getMergedZipPath(outputDir),
CHARACTER_DICTIONARY_MERGED_TITLE,
description,
@@ -633,7 +638,7 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
resolvedMedia.title,
waitForAniListRequestSlot,
);
const storedSnapshot = readSnapshot(getSnapshotPath(outputDir, resolvedMedia.id));
const storedSnapshot = await readSnapshot(getSnapshotPath(outputDir, resolvedMedia.id));
if (!storedSnapshot) {
throw new Error(`Snapshot missing after generation for AniList ${resolvedMedia.id}.`);
}
@@ -642,7 +647,7 @@ export function createCharacterDictionaryRuntimeService(deps: CharacterDictionar
const description = `Character names from ${storedSnapshot.mediaTitle} [AniList media ID ${resolvedMedia.id}]`;
const zipPath = path.join(outputDir, `anilist-${resolvedMedia.id}.zip`);
deps.logInfo?.(`[dictionary] building ZIP for AniList ${resolvedMedia.id}`);
buildDictionaryZip(
await buildDictionaryZip(
zipPath,
dictionaryTitle,
description,
@@ -3,6 +3,7 @@ import * as fs from 'fs';
import * as os from 'os';
import * as path from 'path';
import test from 'node:test';
import { isDeepStrictEqual } from 'node:util';
import { getSnapshotPath, readSnapshot, writeSnapshot } from './cache';
import { CHARACTER_DICTIONARY_FORMAT_VERSION } from './constants';
import type { CharacterDictionarySnapshot } from './types';
@@ -29,17 +30,72 @@ function createSnapshot(): CharacterDictionarySnapshot {
};
}
test('writeSnapshot persists and readSnapshot restores current-format snapshots', () => {
test('writeSnapshot persists and readSnapshot restores current-format snapshots', async () => {
const outputDir = makeTempDir();
const snapshotPath = getSnapshotPath(outputDir, 130298);
const snapshot = createSnapshot();
writeSnapshot(snapshotPath, snapshot);
await writeSnapshot(snapshotPath, snapshot);
assert.deepEqual(readSnapshot(snapshotPath), { ...snapshot, nameSplitSource: 'heuristic' });
assert.deepEqual(await readSnapshot(snapshotPath), { ...snapshot, nameSplitSource: 'heuristic' });
});
test('readSnapshot preserves the mecab name-split source and defaults missing values to heuristic', () => {
// A manual generate and an auto-sync can both land on the same media, so two writes for one
// snapshot can overlap. They must not stream into a shared temp file and interleave into a
// half-and-half snapshot.
test('concurrent writeSnapshot calls for the same media leave one complete snapshot', async () => {
const outputDir = makeTempDir();
const snapshotPath = getSnapshotPath(outputDir, 130298);
const base = createSnapshot();
// Distinct titles, lengths, and term text so the surviving file can be pinned to exactly one
// writer rather than merely "a snapshot that parses". A shared temp file is caught by the
// losing writers failing to rename; interleaved content is only caught when the timing happens
// to leave a mix, which is why the assertion checks identity rather than shape.
const variants: CharacterDictionarySnapshot[] = ['alpha', 'beta', 'gamma'].map((label, index) => {
const entryCount = 400 + index * 100;
return {
...base,
mediaTitle: `${base.mediaTitle} ${label}`,
entryCount,
termEntries: Array.from({ length: entryCount }, (_entry, entryIndex) => [
`${label}${entryIndex}`,
'なまえ',
'name primary',
'',
75,
[`${label} character ${entryIndex} `.repeat(600)],
0,
'',
]) as CharacterDictionarySnapshot['termEntries'],
};
});
await Promise.all(variants.map((variant) => writeSnapshot(snapshotPath, variant)));
const restored = await readSnapshot(snapshotPath);
const expected = variants.map((variant) => ({
...variant,
nameSplitSource: 'heuristic' as const,
}));
const matches = expected.filter((candidate) => isDeepStrictEqual(restored, candidate));
assert.equal(
matches.length,
1,
`expected exactly one writer's complete snapshot to survive, got ${
restored === null
? 'an unreadable file'
: `entryCount=${restored.entryCount}, terms=${restored.termEntries.length}, title=${restored.mediaTitle}`
}`,
);
// Every writer cleaned up after itself, so no temp files are left behind.
const leftovers = fs
.readdirSync(path.dirname(snapshotPath))
.filter((name) => name.includes('.tmp-'));
assert.deepEqual(leftovers, []);
});
test('readSnapshot preserves the mecab name-split source and defaults missing values to heuristic', async () => {
const outputDir = makeTempDir();
const snapshotPath = getSnapshotPath(outputDir, 130298);
const snapshot: CharacterDictionarySnapshot = {
@@ -47,12 +103,12 @@ test('readSnapshot preserves the mecab name-split source and defaults missing va
nameSplitSource: 'mecab',
};
writeSnapshot(snapshotPath, snapshot);
await writeSnapshot(snapshotPath, snapshot);
assert.equal(readSnapshot(snapshotPath)?.nameSplitSource, 'mecab');
assert.equal((await readSnapshot(snapshotPath))?.nameSplitSource, 'mecab');
});
test('readSnapshot ignores snapshots written with an older format version', () => {
test('readSnapshot ignores snapshots written with an older format version', async () => {
const outputDir = makeTempDir();
const snapshotPath = getSnapshotPath(outputDir, 130298);
const staleSnapshot = {
@@ -63,10 +119,10 @@ test('readSnapshot ignores snapshots written with an older format version', () =
fs.mkdirSync(path.dirname(snapshotPath), { recursive: true });
fs.writeFileSync(snapshotPath, JSON.stringify(staleSnapshot), 'utf8');
assert.equal(readSnapshot(snapshotPath), null);
assert.equal(await readSnapshot(snapshotPath), null);
});
test('readSnapshot ignores v15 snapshots with stale romanized character-name entries', () => {
test('readSnapshot ignores v15 snapshots with stale romanized character-name entries', async () => {
const outputDir = makeTempDir();
const snapshotPath = getSnapshotPath(outputDir, 130298);
const staleSnapshot = {
@@ -78,5 +134,5 @@ test('readSnapshot ignores v15 snapshots with stale romanized character-name ent
fs.mkdirSync(path.dirname(snapshotPath), { recursive: true });
fs.writeFileSync(snapshotPath, JSON.stringify(staleSnapshot), 'utf8');
assert.equal(readSnapshot(snapshotPath), null);
assert.equal(await readSnapshot(snapshotPath), null);
});
+83 -10
View File
@@ -102,24 +102,42 @@ export function writeCachedMediaResolution(
writeMediaResolutionEntries(outputDir, [...remaining, normalized]);
}
export function readCachedSnapshots(outputDir: string): CharacterDictionarySnapshot[] {
/**
* Snapshots for long series run to hundreds of MB each, so everything here reads them off the main
* thread's critical path: file IO is async and only the unavoidable JSON.parse runs on the loop,
* one file at a time. Reading the whole directory synchronously used to block the process for
* multiple seconds, long enough for the compositor to declare the app unresponsive mid-playback.
*/
export async function readCachedSnapshots(
outputDir: string,
): Promise<CharacterDictionarySnapshot[]> {
let entries: fs.Dirent[] = [];
try {
entries = fs.readdirSync(getSnapshotsDir(outputDir), { withFileTypes: true });
entries = await fs.promises.readdir(getSnapshotsDir(outputDir), { withFileTypes: true });
} catch {
return [];
}
return entries
const names = entries
.filter((entry) => entry.isFile() && /^anilist-\d+\.json$/.test(entry.name))
.sort((left, right) => left.name.localeCompare(right.name))
.map((entry) => readSnapshot(path.join(getSnapshotsDir(outputDir), entry.name)))
.filter((snapshot): snapshot is CharacterDictionarySnapshot => snapshot !== null);
.map((entry) => entry.name)
.sort((left, right) => left.localeCompare(right));
const snapshots: CharacterDictionarySnapshot[] = [];
for (const name of names) {
const snapshot = await readSnapshot(path.join(getSnapshotsDir(outputDir), name));
if (snapshot) {
snapshots.push(snapshot);
}
}
return snapshots;
}
export function readSnapshot(snapshotPath: string): CharacterDictionarySnapshot | null {
export async function readSnapshot(
snapshotPath: string,
): Promise<CharacterDictionarySnapshot | null> {
try {
const raw = fs.readFileSync(snapshotPath, 'utf8');
const raw = await fs.promises.readFile(snapshotPath, 'utf8');
const parsed = JSON.parse(raw) as Partial<CharacterDictionarySnapshot>;
if (!parsed || typeof parsed !== 'object') {
return null;
@@ -150,9 +168,64 @@ export function readSnapshot(snapshotPath: string): CharacterDictionarySnapshot
}
}
export function writeSnapshot(snapshotPath: string, snapshot: CharacterDictionarySnapshot): void {
// Flushing in a few-MB batches keeps each stringify-and-write slice short; a single
// JSON.stringify of a large snapshot blocks the event loop for seconds.
const SNAPSHOT_WRITE_FLUSH_BYTES = 4 * 1024 * 1024;
// Distinguishes concurrent writes of the same snapshot within one process; the pid alone only
// separates processes, so two overlapping writers would otherwise stream into the same temp file.
let snapshotWriteSequence = 0;
/**
* Streams the snapshot to disk piece by piece instead of stringifying it in one shot, then renames
* the finished file into place so a crash mid-write (or two concurrent writers for the same media)
* can never leave a torn file where a snapshot used to be.
*/
export async function writeSnapshot(
snapshotPath: string,
snapshot: CharacterDictionarySnapshot,
): Promise<void> {
ensureDir(path.dirname(snapshotPath));
fs.writeFileSync(snapshotPath, JSON.stringify(snapshot, null, 2), 'utf8');
snapshotWriteSequence += 1;
const tempPath = `${snapshotPath}.tmp-${process.pid}-${snapshotWriteSequence}`;
const handle = await fs.promises.open(tempPath, 'w');
try {
let buffered: string[] = [];
let bufferedBytes = 0;
const push = async (chunk: string): Promise<void> => {
buffered.push(chunk);
bufferedBytes += chunk.length;
if (bufferedBytes >= SNAPSHOT_WRITE_FLUSH_BYTES) {
const joined = buffered.join('');
buffered = [];
bufferedBytes = 0;
await handle.write(joined, null, 'utf8');
}
};
const writeArray = async (key: string, items: readonly unknown[]): Promise<void> => {
await push(`,${JSON.stringify(key)}:[`);
for (let i = 0; i < items.length; i += 1) {
await push(`${i > 0 ? ',' : ''}${JSON.stringify(items[i])}`);
}
await push(']');
};
const { termEntries, images, ...scalars } = snapshot;
const head = JSON.stringify(scalars);
await push(head.slice(0, -1));
await writeArray('termEntries', termEntries);
await writeArray('images', images);
await push('}');
if (buffered.length > 0) {
await handle.write(buffered.join(''), null, 'utf8');
}
} catch (error) {
await handle.close();
await fs.promises.rm(tempPath, { force: true });
throw error;
}
await handle.close();
await fs.promises.rename(tempPath, snapshotPath);
}
export function buildMergedRevision(
@@ -11,6 +11,22 @@ import {
} from './image-lookup';
import type { CharacterDictionarySnapshot } from './types';
// Lookup indexes rebuild in the background while gets serve stale data, so tests poll until the
// refresh they triggered has landed.
async function waitForRefresh<T>(probe: () => T | null | undefined): Promise<T> {
const deadline = Date.now() + 5000;
for (;;) {
const value = probe();
if (value !== null && value !== undefined) {
return value;
}
if (Date.now() > deadline) {
throw new Error('timed out waiting for background snapshot refresh');
}
await new Promise((resolve) => setTimeout(resolve, 5));
}
}
const PNG_1X1_BASE64 =
'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAwMCAO+nmX8AAAAASUVORK5CYII=';
@@ -18,7 +34,7 @@ function makeTempDir(): string {
return fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-character-image-lookup-'));
}
test('buildCharacterNameImageIndexFromSnapshots maps name terms to character portrait data URLs', () => {
test('buildCharacterNameImageIndexFromSnapshots maps name terms to character portrait data URLs', async () => {
const outputDir = makeTempDir();
const snapshot: CharacterDictionarySnapshot = {
formatVersion: CHARACTER_DICTIONARY_FORMAT_VERSION,
@@ -75,9 +91,9 @@ test('buildCharacterNameImageIndexFromSnapshots maps name terms to character por
{ path: 'img/m130298-va456.png', dataBase64: 'BBBB' },
],
};
writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
await writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
const index = buildCharacterNameImageIndexFromSnapshots(outputDir);
const index = await buildCharacterNameImageIndexFromSnapshots(outputDir);
assert.deepEqual(index.get('アレクシア'), {
src: 'data:image/png;base64,AAAA',
@@ -85,7 +101,7 @@ test('buildCharacterNameImageIndexFromSnapshots maps name terms to character por
});
});
test('buildCharacterNameImageIndexFromSnapshots sniffs image MIME from bytes before path extension', () => {
test('buildCharacterNameImageIndexFromSnapshots sniffs image MIME from bytes before path extension', async () => {
const outputDir = makeTempDir();
const snapshot: CharacterDictionarySnapshot = {
formatVersion: CHARACTER_DICTIONARY_FORMAT_VERSION,
@@ -116,14 +132,14 @@ test('buildCharacterNameImageIndexFromSnapshots sniffs image MIME from bytes bef
],
images: [{ path: 'img/m130298-c123.jpg', dataBase64: PNG_1X1_BASE64 }],
};
writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
await writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
const index = buildCharacterNameImageIndexFromSnapshots(outputDir);
const index = await buildCharacterNameImageIndexFromSnapshots(outputDir);
assert.equal(index.get('アレクシア')?.src, `data:image/png;base64,${PNG_1X1_BASE64}`);
});
test('createCharacterDictionaryImageLookup can scope duplicate names to the current media', () => {
test('createCharacterDictionaryImageLookup can scope duplicate names to the current media', async () => {
const outputDir = makeTempDir();
const towerSnapshot: CharacterDictionarySnapshot = {
formatVersion: CHARACTER_DICTIONARY_FORMAT_VERSION,
@@ -173,15 +189,16 @@ test('createCharacterDictionaryImageLookup can scope duplicate names to the curr
],
images: [{ path: 'img/m21202-c2.png', dataBase64: 'KONOSUBA' }],
};
writeSnapshot(getSnapshotPath(outputDir, towerSnapshot.mediaId), towerSnapshot);
writeSnapshot(getSnapshotPath(outputDir, konosubaSnapshot.mediaId), konosubaSnapshot);
await writeSnapshot(getSnapshotPath(outputDir, towerSnapshot.mediaId), towerSnapshot);
await writeSnapshot(getSnapshotPath(outputDir, konosubaSnapshot.mediaId), konosubaSnapshot);
const lookup = createCharacterDictionaryImageLookup({ outputDir });
assert.equal(lookup.get('カズ', 21202)?.alt, 'Kazuma');
const scoped = await waitForRefresh(() => lookup.get('カズ', 21202));
assert.equal(scoped.alt, 'Kazuma');
});
test('createCharacterDictionaryImageLookup does not fall back globally on scoped miss', () => {
test('createCharacterDictionaryImageLookup does not fall back globally on scoped miss', async () => {
const outputDir = makeTempDir();
const snapshot: CharacterDictionarySnapshot = {
formatVersion: CHARACTER_DICTIONARY_FORMAT_VERSION,
@@ -208,10 +225,11 @@ test('createCharacterDictionaryImageLookup does not fall back globally on scoped
],
images: [{ path: 'img/m115230-c1.png', dataBase64: 'TOWER' }],
};
writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
await writeSnapshot(getSnapshotPath(outputDir, snapshot.mediaId), snapshot);
const lookup = createCharacterDictionaryImageLookup({ outputDir });
const unscoped = await waitForRefresh(() => lookup.get('カズ'));
assert.equal(unscoped.alt, 'Kaz');
assert.equal(lookup.get('カズ', 21202), null);
assert.equal(lookup.get('カズ')?.alt, 'Kaz');
});
@@ -204,11 +204,11 @@ function getSnapshotDirectorySignature(outputDir: string): string {
return parts.sort().join('|');
}
export function buildCharacterNameImageIndexFromSnapshots(
export async function buildCharacterNameImageIndexFromSnapshots(
outputDir: string,
): Map<string, CharacterNameImage> {
): Promise<Map<string, CharacterNameImage>> {
const index = new Map<string, CharacterNameImage>();
for (const snapshot of readCachedSnapshots(outputDir)) {
for (const snapshot of await readCachedSnapshots(outputDir)) {
appendSnapshotImages(index, snapshot);
}
return index;
@@ -228,7 +228,12 @@ export function createCharacterDictionaryImageLookup(deps: {
let signature: string | null = null;
let index = new Map<string, CharacterNameImage>();
let indexByMediaId = new Map<number, Map<string, CharacterNameImage>>();
let refreshInFlight = false;
// Rebuilding means re-reading every cached snapshot (potentially GBs of JSON), which used to run
// synchronously inside a lookup and froze the whole app right after a snapshot changed. Lookups
// now serve the previous index while a single background rebuild catches up; the swap is atomic
// and the signature only advances once the rebuild it belongs to has landed.
function refreshIfNeeded(): void {
if (!outputDir) {
index = new Map<string, CharacterNameImage>();
@@ -237,20 +242,30 @@ export function createCharacterDictionaryImageLookup(deps: {
return;
}
const nextSignature = getSnapshotDirectorySignature(outputDir);
if (nextSignature === signature) {
if (nextSignature === signature || refreshInFlight) {
return;
}
signature = nextSignature;
index = new Map<string, CharacterNameImage>();
indexByMediaId = new Map<number, Map<string, CharacterNameImage>>();
for (const snapshot of readCachedSnapshots(outputDir)) {
appendSnapshotImages(index, snapshot);
const mediaIndex = new Map<string, CharacterNameImage>();
appendSnapshotImages(mediaIndex, snapshot);
if (mediaIndex.size > 0) {
indexByMediaId.set(snapshot.mediaId, mediaIndex);
refreshInFlight = true;
void (async () => {
try {
const snapshots = await readCachedSnapshots(outputDir);
const nextIndex = new Map<string, CharacterNameImage>();
const nextIndexByMediaId = new Map<number, Map<string, CharacterNameImage>>();
for (const snapshot of snapshots) {
appendSnapshotImages(nextIndex, snapshot);
const mediaIndex = new Map<string, CharacterNameImage>();
appendSnapshotImages(mediaIndex, snapshot);
if (mediaIndex.size > 0) {
nextIndexByMediaId.set(snapshot.mediaId, mediaIndex);
}
}
index = nextIndex;
indexByMediaId = nextIndexByMediaId;
signature = nextSignature;
} finally {
refreshInFlight = false;
}
}
})();
}
return {
@@ -32,17 +32,33 @@ function writeSnapshot(outputDir: string, mediaId: number, entries: Array<[strin
);
}
function withTempDir<T>(run: (dir: string) => T): T {
async function withTempDir<T>(run: (dir: string) => Promise<T> | T): Promise<T> {
const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-name-candidates-'));
try {
return run(dir);
return await run(dir);
} finally {
fs.rmSync(dir, { recursive: true, force: true });
}
}
test('collects terms and readings for the current media', () => {
withTempDir((dir) => {
// The snapshot index rebuilds in the background while lookups serve stale data, so tests poll the
// probe until the refresh they triggered has landed.
async function waitForRefresh<T>(probe: () => T | null | undefined): Promise<T> {
const deadline = Date.now() + 5000;
for (;;) {
const value = probe();
if (value !== null && value !== undefined) {
return value;
}
if (Date.now() > deadline) {
throw new Error('timed out waiting for background snapshot refresh');
}
await new Promise((resolve) => setTimeout(resolve, 5));
}
}
test('collects terms and readings for the current media', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [
['ミナト', 'みなと'],
['湊', 'みなと'],
@@ -53,17 +69,16 @@ test('collects terms and readings for the current media', () => {
outputDir: dir,
getCurrentMediaId: () => 1,
});
const candidates = lookup.get();
const candidates = await waitForRefresh(() => lookup.get());
assert.ok(candidates);
assert.deepEqual([...candidates.forms].sort(), ['みなと', 'ミナト', '湊'].sort());
// Deduplicated: both entries share the みなと reading.
assert.equal(candidates.forms.length, 3);
});
});
test('returns null without a media scope so the scanner stays exhaustive', () => {
withTempDir((dir) => {
test('returns null without a media scope so the scanner stays exhaustive', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [['ミナト', 'みなと']]);
const lookup = createCharacterNameCandidateLookup({
@@ -71,12 +86,14 @@ test('returns null without a media scope so the scanner stays exhaustive', () =>
getCurrentMediaId: () => null,
});
// The explicitly-scoped probe proves the index has loaded before the unscoped case is judged.
await waitForRefresh(() => lookup.get(1));
assert.equal(lookup.get(), null);
});
});
test('returns null for a media with no cached snapshot', () => {
withTempDir((dir) => {
test('returns null for a media with no cached snapshot', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [['ミナト', 'みなと']]);
const lookup = createCharacterNameCandidateLookup({
@@ -84,29 +101,31 @@ test('returns null for a media with no cached snapshot', () => {
getCurrentMediaId: () => 999,
});
await waitForRefresh(() => lookup.get(1));
assert.equal(lookup.get(), null);
});
});
test('key changes when the snapshot content changes', () => {
withTempDir((dir) => {
test('key changes when the snapshot content changes', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [['ミナト', 'みなと']]);
const lookup = createCharacterNameCandidateLookup({
outputDir: dir,
getCurrentMediaId: () => 1,
});
const first = lookup.get();
const first = await waitForRefresh(() => lookup.get());
writeSnapshot(dir, 1, [
['ミナト', 'みなと'],
['アクア', 'あくあ'],
]);
lookup.invalidate();
const second = lookup.get();
const second = await waitForRefresh(() => {
const candidates = lookup.get();
return candidates && candidates.forms.length === 4 ? candidates : null;
});
assert.ok(first && second);
assert.notEqual(first.key, second.key);
assert.equal(second.forms.length, 4);
});
});
@@ -114,8 +133,8 @@ test('key changes when the snapshot content changes', () => {
// directory every call. Asserted behaviorally: an unannounced on-disk change is
// invisible until the recheck interval elapses, which can only be true if the
// filesystem is not consulted per lookup.
test('does not re-read the snapshot directory on every lookup', () => {
withTempDir((dir) => {
test('does not re-read the snapshot directory on every lookup', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [['ミナト', 'みなと']]);
let nowMs = 1_000_000;
const lookup = createCharacterNameCandidateLookup({
@@ -124,6 +143,7 @@ test('does not re-read the snapshot directory on every lookup', () => {
now: () => nowMs,
});
await waitForRefresh(() => lookup.get());
assert.equal(lookup.get()?.forms.length, 2);
writeSnapshot(dir, 1, [
@@ -135,12 +155,16 @@ test('does not re-read the snapshot directory on every lookup', () => {
assert.equal(lookup.get()?.forms.length, 2, 'expected the cached list within the interval');
nowMs += 10_000;
assert.equal(lookup.get()?.forms.length, 4, 'expected a refresh past the interval');
const refreshed = await waitForRefresh(() => {
const candidates = lookup.get();
return candidates && candidates.forms.length === 4 ? candidates : null;
});
assert.equal(refreshed.forms.length, 4, 'expected a refresh past the interval');
});
});
test('invalidate picks up a snapshot change immediately', () => {
withTempDir((dir) => {
test('invalidate picks up a snapshot change on the next refresh', async () => {
await withTempDir(async (dir) => {
writeSnapshot(dir, 1, [['ミナト', 'みなと']]);
let nowMs = 1_000_000;
const lookup = createCharacterNameCandidateLookup({
@@ -149,6 +173,7 @@ test('invalidate picks up a snapshot change immediately', () => {
now: () => nowMs,
});
await waitForRefresh(() => lookup.get());
assert.equal(lookup.get()?.forms.length, 2);
writeSnapshot(dir, 1, [
@@ -158,6 +183,10 @@ test('invalidate picks up a snapshot change immediately', () => {
nowMs += 1;
lookup.invalidate();
assert.equal(lookup.get()?.forms.length, 4);
const refreshed = await waitForRefresh(() => {
const candidates = lookup.get();
return candidates && candidates.forms.length === 4 ? candidates : null;
});
assert.equal(refreshed.forms.length, 4);
});
});
@@ -98,7 +98,12 @@ export function createCharacterNameCandidateLookup(deps: {
let signature: string | null = null;
let lastSignatureCheckAtMs = 0;
let formsByMediaId = new Map<number, string[]>();
let refreshInFlight = false;
// Same stale-while-revalidate shape as the image lookup: the rebuild re-reads every cached
// snapshot, so it runs in the background while lookups keep serving the previous forms. The
// signature only advances once its rebuild has landed, so a failed or superseded rebuild is
// retried on the next signature check.
function refreshIfNeeded(): void {
if (!outputDir) {
formsByMediaId = new Map<number, string[]>();
@@ -114,17 +119,26 @@ export function createCharacterNameCandidateLookup(deps: {
}
lastSignatureCheckAtMs = nowMs;
const nextSignature = getSnapshotDirectorySignature(outputDir);
if (nextSignature === signature) {
if (nextSignature === signature || refreshInFlight) {
return;
}
signature = nextSignature;
formsByMediaId = new Map<number, string[]>();
for (const snapshot of readCachedSnapshots(outputDir)) {
const forms = collectSnapshotNameForms(snapshot);
if (forms.length > 0) {
formsByMediaId.set(snapshot.mediaId, forms);
refreshInFlight = true;
void (async () => {
try {
const snapshots = await readCachedSnapshots(outputDir);
const nextFormsByMediaId = new Map<number, string[]>();
for (const snapshot of snapshots) {
const forms = collectSnapshotNameForms(snapshot);
if (forms.length > 0) {
nextFormsByMediaId.set(snapshot.mediaId, forms);
}
}
formsByMediaId = nextFormsByMediaId;
signature = nextSignature;
} finally {
refreshInFlight = false;
}
}
})();
}
return {
@@ -34,7 +34,7 @@ function createSnapshotWithoutImages(): CharacterDictionarySnapshot {
test('generateForCurrentMedia refreshes same-version snapshots missing images when inline images are enabled', async () => {
const userDataPath = makeTempDir();
const outputDir = path.join(userDataPath, 'character-dictionaries');
writeSnapshot(getSnapshotPath(outputDir, 130298), createSnapshotWithoutImages());
await writeSnapshot(getSnapshotPath(outputDir, 130298), createSnapshotWithoutImages());
const originalFetch = globalThis.fetch;
const fetchUrls: string[] = [];
@@ -124,7 +124,7 @@ test('generateForCurrentMedia refreshes same-version snapshots missing images wh
test('generateForCurrentMedia keeps failed MeCab name split refreshes retryable', async () => {
const userDataPath = makeTempDir();
const outputDir = path.join(userDataPath, 'character-dictionaries');
writeSnapshot(getSnapshotPath(outputDir, 130298), {
await writeSnapshot(getSnapshotPath(outputDir, 130298), {
...createSnapshotWithoutImages(),
nameSplitSource: 'heuristic',
});
@@ -213,7 +213,7 @@ test('generateForCurrentMedia keeps failed MeCab name split refreshes retryable'
test('generateForCurrentMedia keeps mecab-split snapshots when MeCab is available', async () => {
const userDataPath = makeTempDir();
const outputDir = path.join(userDataPath, 'character-dictionaries');
writeSnapshot(getSnapshotPath(outputDir, 130298), {
await writeSnapshot(getSnapshotPath(outputDir, 130298), {
...createSnapshotWithoutImages(),
nameSplitSource: 'mecab',
});
@@ -253,7 +253,7 @@ test('generateForCurrentMedia keeps mecab-split snapshots when MeCab is availabl
test('generateForCurrentMedia keeps heuristic-split snapshots while MeCab is unavailable', async () => {
const userDataPath = makeTempDir();
const outputDir = path.join(userDataPath, 'character-dictionaries');
writeSnapshot(getSnapshotPath(outputDir, 130298), {
await writeSnapshot(getSnapshotPath(outputDir, 130298), {
...createSnapshotWithoutImages(),
nameSplitSource: 'heuristic',
});
@@ -293,7 +293,7 @@ test('generateForCurrentMedia keeps heuristic-split snapshots while MeCab is una
test('generateForCurrentMedia keeps same-version snapshots without images when inline images are disabled', async () => {
const userDataPath = makeTempDir();
const outputDir = path.join(userDataPath, 'character-dictionaries');
writeSnapshot(getSnapshotPath(outputDir, 130298), createSnapshotWithoutImages());
await writeSnapshot(getSnapshotPath(outputDir, 130298), createSnapshotWithoutImages());
const originalFetch = globalThis.fetch;
globalThis.fetch = (async (input: string | URL | Request) => {
@@ -42,7 +42,7 @@ function readStoredZipEntries(zipPath: string): Map<string, Buffer> {
return entries;
}
test('buildDictionaryZip writes a valid stored zip without fs.writeFileSync', () => {
test('buildDictionaryZip writes a valid stored zip without fs.writeFileSync', async () => {
const tempDir = makeTempDir();
const outputPath = path.join(tempDir, 'dictionary.zip');
const termEntries: CharacterDictionaryTermEntry[] = [
@@ -62,7 +62,7 @@ test('buildDictionaryZip writes a valid stored zip without fs.writeFileSync', ()
);
}) as typeof Buffer.concat;
const result = buildDictionaryZip(
const result = await buildDictionaryZip(
outputPath,
'Dictionary Title',
'Dictionary Description',
@@ -106,11 +106,11 @@ test('buildDictionaryZip writes a valid stored zip without fs.writeFileSync', ()
}
});
test('readDictionaryZipRevision reads the built revision and rejects foreign archives', () => {
test('readDictionaryZipRevision reads the built revision and rejects foreign archives', async () => {
const dir = makeTempDir();
try {
const zipPath = path.join(dir, 'merged.zip');
buildDictionaryZip(
await buildDictionaryZip(
zipPath,
'SubMiner Character Dictionary',
'Character names',
+9 -5
View File
@@ -1,5 +1,5 @@
import * as path from 'path';
import { readStoredZipFirstFile, writeStoredZip } from '../../shared/stored-zip';
import { readStoredZipFirstFile, writeStoredZipAsync } from '../../shared/stored-zip';
import { ensureDir } from './fs-utils';
import type { CharacterDictionarySnapshotImage, CharacterDictionaryTermEntry } from './types';
@@ -48,14 +48,14 @@ export function readDictionaryZipRevision(zipPath: string): string | null {
}
}
export function buildDictionaryZip(
export async function buildDictionaryZip(
outputPath: string,
dictionaryTitle: string,
description: string,
revision: string,
termEntries: CharacterDictionaryTermEntry[],
images: CharacterDictionarySnapshotImage[],
): { zipPath: string; entryCount: number } {
): Promise<{ zipPath: string; entryCount: number }> {
ensureDir(path.dirname(outputPath));
function* zipFiles(): Iterable<{ name: string; data: Buffer }> {
@@ -78,7 +78,11 @@ export function buildDictionaryZip(
};
}
const entriesPerBank = 10_000;
// Each bank is stringified in one shot, so the bank size sets the longest single block in the
// build. 10k entries measured ~38MB and ~135ms per bank on a real merged dictionary; 2k keeps
// every bank under the archive writer's yield budget at ~27ms. Yomitan reads any number of
// term_bank_N.json files, so this only changes how the terms are split across them.
const entriesPerBank = 2_000;
for (let i = 0; i < termEntries.length; i += entriesPerBank) {
yield {
name: `term_bank_${Math.floor(i / entriesPerBank) + 1}.json`,
@@ -87,6 +91,6 @@ export function buildDictionaryZip(
}
}
writeStoredZip(outputPath, zipFiles());
await writeStoredZipAsync(outputPath, zipFiles());
return { zipPath: outputPath, entryCount: termEntries.length };
}
+74
View File
@@ -0,0 +1,74 @@
import assert from 'node:assert/strict';
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import test from 'node:test';
import { readStoredZipFirstFile, writeStoredZip, writeStoredZipAsync } from './stored-zip';
function makeTempDir(): string {
return fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-stored-zip-'));
}
function readEntries(zipPath: string): Map<string, Buffer> {
const archive = fs.readFileSync(zipPath);
const entries = new Map<string, Buffer>();
let cursor = 0;
while (cursor + 4 <= archive.length) {
const signature = archive.readUInt32LE(cursor);
if (signature === 0x02014b50 || signature === 0x06054b50) {
break;
}
assert.equal(signature, 0x04034b50, `unexpected local file header at offset ${cursor}`);
const size = archive.readUInt32LE(cursor + 18);
const nameLength = archive.readUInt16LE(cursor + 26);
const extraLength = archive.readUInt16LE(cursor + 28);
const nameStart = cursor + 30;
const dataStart = nameStart + nameLength + extraLength;
entries.set(
archive.subarray(nameStart, nameStart + nameLength).toString('utf8'),
Buffer.from(archive.subarray(dataStart, dataStart + size)),
);
cursor = dataStart + size;
}
return entries;
}
// The async writer yields on a byte budget between entries, so an entry that is itself larger than
// that budget is the case where the accounting could drift: the offsets, CRCs, and central
// directory all have to come out identical to the synchronous writer.
test('writeStoredZipAsync writes a correct archive when one entry exceeds the yield budget', async () => {
const dir = makeTempDir();
try {
// Comfortably past the writer's 8MB yield budget.
const oversized = Buffer.alloc(10 * 1024 * 1024);
for (let i = 0; i < oversized.length; i += 1) {
oversized[i] = i % 251;
}
const files = [
{ name: 'index.json', data: Buffer.from('{"revision":"rev-1"}', 'utf8') },
{ name: 'big.bin', data: oversized },
{ name: 'after.txt', data: Buffer.from('written after the oversized entry', 'utf8') },
];
const asyncPath = path.join(dir, 'async.zip');
const syncPath = path.join(dir, 'sync.zip');
const asyncResult = await writeStoredZipAsync(asyncPath, files);
const syncResult = writeStoredZip(syncPath, files);
assert.equal(asyncResult.entryCount, 3);
assert.deepEqual(asyncResult, syncResult);
// Byte-identical to the synchronous writer: yielding mid-archive changed no offset or CRC.
assert.ok(fs.readFileSync(asyncPath).equals(fs.readFileSync(syncPath)));
const entries = readEntries(asyncPath);
assert.deepEqual([...entries.keys()], ['index.json', 'big.bin', 'after.txt']);
assert.ok(entries.get('big.bin')!.equals(oversized));
assert.equal(entries.get('after.txt')!.toString('utf8'), 'written after the oversized entry');
// The trailing records still parse, which is what proves the archive is complete.
assert.equal(readStoredZipFirstFile(asyncPath)?.name, 'index.json');
} finally {
fs.rmSync(dir, { recursive: true, force: true });
}
});
+103 -47
View File
@@ -1,4 +1,5 @@
import * as fs from 'fs';
import * as zlib from 'zlib';
type ZipEntry = {
name: string;
@@ -36,10 +37,17 @@ const CRC32_TABLE = (() => {
return table;
})();
// Native CRC32 (Node >= 20.15) runs at native throughput, which matters for the multi-hundred-MB
// dictionary archives; the table loop stays as a fallback for runtimes without it.
const nativeCrc32 = (zlib as { crc32?: (data: Uint8Array, value?: number) => number }).crc32;
function crc32(data: Buffer): number {
if (typeof nativeCrc32 === 'function') {
return nativeCrc32(data) >>> 0;
}
let crc = 0xffffffff;
for (const byte of data) {
crc = CRC32_TABLE[(crc ^ byte) & 0xff]! ^ (crc >>> 8);
for (let i = 0; i < data.length; i += 1) {
crc = CRC32_TABLE[(crc ^ data[i]!) & 0xff]! ^ (crc >>> 8);
}
return (crc ^ 0xffffffff) >>> 0;
}
@@ -294,59 +302,70 @@ function writeBuffer(fd: number, buffer: Buffer): void {
}
}
type ZipWriteState = {
entries: ZipEntry[];
offset: number;
};
/** Appends one stored entry (local header + data) and returns the bytes written. */
function appendStoredZipFile(fd: number, state: ZipWriteState, file: StoredZipFile): number {
const fileName = Buffer.from(file.name, 'utf8');
const fileSize = file.data.length;
if (fileName.length > ZIP32_MAX_UINT16) {
throw new RangeError(`ZIP entry name too long: ${file.name}`);
}
if (fileSize > ZIP32_MAX_UINT32) {
throw new RangeError(`ZIP entry too large for ZIP32: ${file.name}`);
}
if (state.offset > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
const fileCrc32 = crc32(file.data);
const localHeader = createLocalFileHeader(fileName, fileCrc32, fileSize);
const nextOffset = state.offset + localHeader.length + fileSize;
if (nextOffset > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
writeBuffer(fd, localHeader);
writeBuffer(fd, file.data);
state.entries.push({
name: file.name,
crc32: fileCrc32,
size: fileSize,
localHeaderOffset: state.offset,
});
const written = nextOffset - state.offset;
state.offset = nextOffset;
return written;
}
function finishStoredZip(fd: number, state: ZipWriteState): void {
const centralStart = state.offset;
if (centralStart > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
for (const entry of state.entries) {
const centralHeader = createCentralDirectoryHeader(entry);
writeBuffer(fd, centralHeader);
state.offset += centralHeader.length;
}
const centralSize = state.offset - centralStart;
writeBuffer(fd, createEndOfCentralDirectory(state.entries.length, centralSize, centralStart));
}
export function writeStoredZip(
outputPath: string,
files: Iterable<StoredZipFile>,
): { entryCount: number } {
const entries: ZipEntry[] = [];
let offset = 0;
const state: ZipWriteState = { entries: [], offset: 0 };
const fd = fs.openSync(outputPath, 'w');
try {
for (const file of files) {
const fileName = Buffer.from(file.name, 'utf8');
const fileSize = file.data.length;
if (fileName.length > ZIP32_MAX_UINT16) {
throw new RangeError(`ZIP entry name too long: ${file.name}`);
}
if (fileSize > ZIP32_MAX_UINT32) {
throw new RangeError(`ZIP entry too large for ZIP32: ${file.name}`);
}
if (offset > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
const fileCrc32 = crc32(file.data);
const localHeader = createLocalFileHeader(fileName, fileCrc32, fileSize);
const nextOffset = offset + localHeader.length + fileSize;
if (nextOffset > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
writeBuffer(fd, localHeader);
writeBuffer(fd, file.data);
entries.push({
name: file.name,
crc32: fileCrc32,
size: fileSize,
localHeaderOffset: offset,
});
if (nextOffset > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
offset = nextOffset;
appendStoredZipFile(fd, state, file);
}
const centralStart = offset;
if (centralStart > ZIP32_MAX_UINT32) {
throw new RangeError('Archive exceeds ZIP32 limits (Zip64 not implemented)');
}
for (const entry of entries) {
const centralHeader = createCentralDirectoryHeader(entry);
writeBuffer(fd, centralHeader);
offset += centralHeader.length;
}
const centralSize = offset - centralStart;
writeBuffer(fd, createEndOfCentralDirectory(entries.length, centralSize, centralStart));
finishStoredZip(fd, state);
} catch (error) {
fs.closeSync(fd);
fs.rmSync(outputPath, { force: true });
@@ -354,5 +373,42 @@ export function writeStoredZip(
}
fs.closeSync(fd);
return { entryCount: entries.length };
return { entryCount: state.entries.length };
}
// Yielding roughly every 8MB keeps individual event-loop blocks in the low tens of milliseconds
// while adding a negligible number of macrotask hops even for the largest merged dictionary.
const ASYNC_ZIP_YIELD_BYTE_BUDGET = 8 * 1024 * 1024;
/**
* Same archive as {@link writeStoredZip}, written without starving the event loop: entry
* generation, CRC, and writes proceed in byte-budgeted slices with a macrotask yield in between.
* Multi-hundred-MB dictionary archives previously blocked the main process long enough for the
* compositor to declare the app unresponsive.
*/
export async function writeStoredZipAsync(
outputPath: string,
files: Iterable<StoredZipFile>,
): Promise<{ entryCount: number }> {
const state: ZipWriteState = { entries: [], offset: 0 };
const fd = fs.openSync(outputPath, 'w');
try {
let bytesSinceYield = 0;
for (const file of files) {
bytesSinceYield += appendStoredZipFile(fd, state, file);
if (bytesSinceYield >= ASYNC_ZIP_YIELD_BYTE_BUDGET) {
bytesSinceYield = 0;
await new Promise<void>((resolve) => setImmediate(resolve));
}
}
finishStoredZip(fd, state);
} catch (error) {
fs.closeSync(fd);
fs.rmSync(outputPath, { force: true });
throw error;
}
fs.closeSync(fd);
return { entryCount: state.entries.length };
}