feat: add asset import services
This commit is contained in:
1 parent
4b596226e0
commit
3497df9dfd
46 files changed
+5383
No files matched your search
@@ -0,0 +1,22 @@
|
||||
import path from 'node:path'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { resolveGamedataFile } from './asset-paths'
|
||||
|
||||
describe('resolveGamedataFile', () => {
|
||||
it('uses absolute external paths as-is', () => {
|
||||
const absolutePath = path.win32.join(
|
||||
'E:\\',
|
||||
'Users',
|
||||
'simol',
|
||||
'Desktop',
|
||||
'DEV',
|
||||
'Nitro-Files',
|
||||
'nitro-assets',
|
||||
'gamedata',
|
||||
'FigureMap.json',
|
||||
)
|
||||
|
||||
expect(resolveGamedataFile(absolutePath, 'FigureMap.json')).toBe(absolutePath)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import { existsSync } from 'node:fs'
|
||||
import path from 'node:path'
|
||||
|
||||
// The CMS default gamedata JSON paths point under public/nitro-assets/gamedata/,
|
||||
// but this deployment keeps the gamedata JSON (FigureMap/FigureData/EffectMap…)
|
||||
// under public/Gamedata/config/. `resolveGamedataFile` falls back there when the
|
||||
// configured path is absent, so importers read AND write the real file — no
|
||||
// manual setting required. (Bundled .nitro stay under public/nitro-assets/bundled,
|
||||
// the client's served asset dir — see *_NITRO_DIR consts in the import services.)
|
||||
|
||||
function pub(...segs: string[]): string {
|
||||
return path.join(/*turbopackIgnore: true*/ process.cwd(), 'public', ...segs)
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a gamedata JSON file from a `*_url` setting value (e.g.
|
||||
* `/nitro-assets/gamedata/FigureMap.json`). Falls back to
|
||||
* `public/Gamedata/config/<fallbackName>` if the configured path is absent.
|
||||
* Returns the configured path when neither exists (write target for fresh setups).
|
||||
*/
|
||||
export function resolveGamedataFile(url: string, fallbackName: string): string {
|
||||
if (path.win32.isAbsolute(url) || (path.isAbsolute(url) && !url.startsWith('/'))) return url
|
||||
const configured = pub(url.replace(/^\/+/, ''))
|
||||
if (existsSync(/*turbopackIgnore: true*/ configured)) return configured
|
||||
const gamedata = pub('Gamedata', 'config', fallbackName)
|
||||
if (existsSync(/*turbopackIgnore: true*/ gamedata)) return gamedata
|
||||
return configured
|
||||
}
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { validatePngBytes, validateSwfBytes } from './download'
|
||||
|
||||
describe('import/core/download validators', () => {
|
||||
it('accepts FWS/CWS/ZWS swf magic', () => {
|
||||
expect(validateSwfBytes(Buffer.from('FWS\x06\x00\x00\x00\x00'))).toBe(true)
|
||||
expect(validateSwfBytes(Buffer.from('CWS\x06\x00\x00\x00\x00'))).toBe(true)
|
||||
expect(validateSwfBytes(Buffer.from('ZWS\x06\x00\x00\x00\x00'))).toBe(true)
|
||||
})
|
||||
|
||||
it('rejects non-swf and too-short buffers', () => {
|
||||
expect(validateSwfBytes(Buffer.from('PNG\x00\x00\x00\x00\x00'))).toBe(false)
|
||||
expect(validateSwfBytes(Buffer.from('FW'))).toBe(false)
|
||||
})
|
||||
|
||||
it('accepts png magic and rejects others', () => {
|
||||
expect(validatePngBytes(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]))).toBe(true)
|
||||
expect(validatePngBytes(Buffer.from([1, 2, 3, 4, 5, 6, 7, 8]))).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
import { promises as fs } from 'node:fs'
|
||||
|
||||
export function validateSwfBytes(buffer: Buffer): boolean {
|
||||
if (buffer.length < 8) return false
|
||||
const sig = buffer.toString('ascii', 0, 3)
|
||||
return sig === 'FWS' || sig === 'CWS' || sig === 'ZWS'
|
||||
}
|
||||
|
||||
export function validatePngBytes(buffer: Buffer): boolean {
|
||||
if (buffer.length < 8) return false
|
||||
return buffer[0] === 137 && buffer[1] === 80 && buffer[2] === 78 && buffer[3] === 71
|
||||
}
|
||||
|
||||
export async function downloadFile(
|
||||
url: string,
|
||||
destPath: string,
|
||||
options?: { maxRetries?: number; validate?: 'swf' | 'png' },
|
||||
): Promise<{ ok: boolean; size: number }> {
|
||||
const maxRetries = options?.maxRetries ?? 3
|
||||
const baseDelay = 1000
|
||||
|
||||
for (let attempt = 0; attempt <= maxRetries; attempt++) {
|
||||
try {
|
||||
if (attempt > 0) {
|
||||
await new Promise((r) => setTimeout(r, baseDelay * 2 ** (attempt - 1)))
|
||||
}
|
||||
const res = await fetch(url, { signal: AbortSignal.timeout(15000) })
|
||||
if (!res.ok) {
|
||||
const deterministic = res.status === 404 || res.status === 410 || res.status === 403
|
||||
if (deterministic || attempt === maxRetries) {
|
||||
console.warn(
|
||||
`[import-download] Download failed ${url}: ${res.status}${deterministic ? ' (deterministic, not retrying)' : ` after ${maxRetries + 1} attempts`}`,
|
||||
)
|
||||
return { ok: false, size: 0 }
|
||||
}
|
||||
continue
|
||||
}
|
||||
const buffer = Buffer.from(await res.arrayBuffer())
|
||||
if (buffer.length < 8) {
|
||||
if (attempt === maxRetries) return { ok: false, size: 0 }
|
||||
continue
|
||||
}
|
||||
if (options?.validate === 'swf' && !validateSwfBytes(buffer)) {
|
||||
if (attempt === maxRetries) {
|
||||
console.warn(`[import-download] Invalid SWF magic bytes from ${url}`)
|
||||
return { ok: false, size: 0 }
|
||||
}
|
||||
continue
|
||||
}
|
||||
if (options?.validate === 'png' && !validatePngBytes(buffer)) {
|
||||
if (attempt === maxRetries) {
|
||||
console.warn(`[import-download] Invalid PNG magic bytes from ${url}`)
|
||||
return { ok: false, size: 0 }
|
||||
}
|
||||
continue
|
||||
}
|
||||
await fs.writeFile(/*turbopackIgnore: true*/ destPath, buffer)
|
||||
return { ok: true, size: buffer.length }
|
||||
} catch (err) {
|
||||
if (attempt === maxRetries) {
|
||||
console.warn(`[import-download] Download error ${url}:`, (err as Error).message)
|
||||
return { ok: false, size: 0 }
|
||||
}
|
||||
}
|
||||
}
|
||||
return { ok: false, size: 0 }
|
||||
}
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
import { promises as fs } from 'node:fs'
|
||||
import os from 'node:os'
|
||||
import path from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
||||
import { readGamedataJson, withGamedataLock, writeGamedataJsonAtomic } from './gamedata-json'
|
||||
|
||||
let dir: string
|
||||
let file: string
|
||||
|
||||
beforeEach(async () => {
|
||||
dir = await fs.mkdtemp(path.join(os.tmpdir(), 'gamedata-'))
|
||||
file = path.join(dir, 'EffectMap.json')
|
||||
await fs.writeFile(file, JSON.stringify({ effects: [{ id: '1', lib: 'A' }] }))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await fs.rm(dir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('import/core/gamedata-json', () => {
|
||||
it('reads parsed JSON', async () => {
|
||||
const data = await readGamedataJson<{ effects: unknown[] }>(file)
|
||||
expect(data.effects).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('writes atomically (round-trips)', async () => {
|
||||
await writeGamedataJsonAtomic(file, { effects: [{ id: '9', lib: 'Z' }] })
|
||||
const data = await readGamedataJson<{ effects: Array<{ id: string }> }>(file)
|
||||
expect(data.effects[0].id).toBe('9')
|
||||
})
|
||||
|
||||
it('serializes concurrent writers (no lost update)', async () => {
|
||||
const bump = (n: number) =>
|
||||
withGamedataLock(file, async () => {
|
||||
const d = await readGamedataJson<{ effects: unknown[] }>(file)
|
||||
d.effects.push({ id: String(n), lib: `L${n}` })
|
||||
await writeGamedataJsonAtomic(file, d)
|
||||
})
|
||||
await Promise.all([bump(2), bump(3), bump(4)])
|
||||
const d = await readGamedataJson<{ effects: unknown[] }>(file)
|
||||
expect(d.effects).toHaveLength(4) // 1 seed + 3 appended, none lost
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
import { promises as fs } from 'node:fs'
|
||||
|
||||
const LOCK_STALE_MS = 60_000
|
||||
const LOCK_ACQUIRE_TIMEOUT_MS = 30_000
|
||||
|
||||
// In-process serialization, keyed by absolute file path, layered on top of the
|
||||
// on-disk lock so multiple awaiters in the same Node process queue cleanly.
|
||||
const chains = new Map<string, Promise<void>>()
|
||||
|
||||
async function acquireDiskLock(lockPath: string): Promise<void> {
|
||||
const deadline = Date.now() + LOCK_ACQUIRE_TIMEOUT_MS
|
||||
while (true) {
|
||||
try {
|
||||
const fh = await fs.open(/*turbopackIgnore: true*/ lockPath, 'wx')
|
||||
await fh.writeFile(String(process.pid))
|
||||
await fh.close()
|
||||
return
|
||||
} catch (err) {
|
||||
const e = err as NodeJS.ErrnoException
|
||||
if (e.code !== 'EEXIST') throw err
|
||||
try {
|
||||
const stat = await fs.stat(/*turbopackIgnore: true*/ lockPath)
|
||||
if (Date.now() - stat.mtimeMs > LOCK_STALE_MS) {
|
||||
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {})
|
||||
continue
|
||||
}
|
||||
} catch {
|
||||
/* lock vanished between checks */
|
||||
}
|
||||
if (Date.now() > deadline) throw new Error(`gamedata lock timeout: ${lockPath}`)
|
||||
await new Promise((r) => setTimeout(r, 100))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function withGamedataLock<T>(filePath: string, fn: () => Promise<T>): Promise<T> {
|
||||
let release!: () => void
|
||||
const acquired = new Promise<void>((r) => {
|
||||
release = r
|
||||
})
|
||||
const prev = chains.get(filePath) ?? Promise.resolve()
|
||||
chains.set(filePath, acquired)
|
||||
await prev
|
||||
const lockPath = `${filePath}.lock`
|
||||
await acquireDiskLock(lockPath)
|
||||
try {
|
||||
return await fn()
|
||||
} finally {
|
||||
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {})
|
||||
release()
|
||||
}
|
||||
}
|
||||
|
||||
export async function readGamedataJson<T>(filePath: string): Promise<T> {
|
||||
const raw = await fs.readFile(/*turbopackIgnore: true*/ filePath, 'utf-8')
|
||||
return JSON.parse(raw) as T
|
||||
}
|
||||
|
||||
export async function writeGamedataJsonAtomic(filePath: string, data: unknown): Promise<void> {
|
||||
const tmp = `${filePath}.tmp`
|
||||
await fs.writeFile(/*turbopackIgnore: true*/ tmp, JSON.stringify(data, null, 2), 'utf-8')
|
||||
await fs.rename(/*turbopackIgnore: true*/ tmp, filePath)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { __clearGordonCache, resolveGordonBuildUrl } from './gordon'
|
||||
|
||||
beforeEach(() => {
|
||||
__clearGordonCache()
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
afterEach(() => vi.restoreAllMocks())
|
||||
|
||||
describe('import/core/gordon', () => {
|
||||
it('parses flash.client.url from external_variables and normalizes', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
new Response('flash.client.url=//images.habbo.com/gordon/PRODUCTION-XYZ/\n', { status: 200 }),
|
||||
)
|
||||
expect(await resolveGordonBuildUrl()).toBe('https://images.habbo.com/gordon/PRODUCTION-XYZ/')
|
||||
})
|
||||
|
||||
it('caches the resolved build (second call does not refetch)', async () => {
|
||||
const spy = vi
|
||||
.spyOn(globalThis, 'fetch')
|
||||
.mockResolvedValue(new Response('flash.client.url=//x/gordon/B/\n', { status: 200 }))
|
||||
await resolveGordonBuildUrl()
|
||||
await resolveGordonBuildUrl()
|
||||
expect(spy).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('throws when flash.client.url is absent', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(new Response('other=1\n', { status: 200 }))
|
||||
await expect(resolveGordonBuildUrl()).rejects.toThrow()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
const EXTERNAL_VARIABLES_URL = 'https://www.habbo.com/gamedata/external_variables/1'
|
||||
const CACHE_TTL = 30 * 60 * 1000
|
||||
|
||||
let gordonCache: { url: string; ts: number } | null = null
|
||||
|
||||
export function __clearGordonCache(): void {
|
||||
gordonCache = null
|
||||
}
|
||||
|
||||
/** Resolve the current Habbo "gordon" asset build URL (trailing slash) from external_variables. */
|
||||
export async function resolveGordonBuildUrl(): Promise<string> {
|
||||
const now = Date.now()
|
||||
if (gordonCache && now - gordonCache.ts < CACHE_TTL) return gordonCache.url
|
||||
const res = await fetch(EXTERNAL_VARIABLES_URL, { signal: AbortSignal.timeout(20000) })
|
||||
if (!res.ok) throw new Error(`external_variables fetch failed: ${res.status}`)
|
||||
const text = await res.text()
|
||||
const line = /^flash\.client\.url=(.+)$/m.exec(text)?.[1]?.trim()
|
||||
if (!line) throw new Error('flash.client.url not found in external_variables')
|
||||
const normalized = (line.startsWith('//') ? `https:${line}` : line).replace(/\/?$/, '/')
|
||||
gordonCache = { url: normalized, ts: now }
|
||||
return normalized
|
||||
}
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { runSseBatch } from './sse-batch'
|
||||
|
||||
async function collect(res: Response): Promise<Array<Record<string, unknown>>> {
|
||||
const text = await res.text()
|
||||
return text
|
||||
.split('\n\n')
|
||||
.filter((l) => l.startsWith('data: '))
|
||||
.map((l) => JSON.parse(l.slice(6)))
|
||||
}
|
||||
|
||||
describe('import/core/sse-batch', () => {
|
||||
it('emits batch_start, per-item, and batch_complete with correct counts', async () => {
|
||||
const res = runSseBatch({
|
||||
items: ['a', 'b', 'c'],
|
||||
concurrency: 2,
|
||||
worker: async (item) => ({ ok: item !== 'b', error: item === 'b' ? 'boom' : undefined }),
|
||||
labelOf: (item) => item,
|
||||
})
|
||||
const events = await collect(res)
|
||||
expect(events[0]).toMatchObject({ type: 'batch_start', total: 3, concurrency: 2 })
|
||||
const done = events.filter((e) => e.type === 'item_progress' && e.status === 'done')
|
||||
const failed = events.filter((e) => e.type === 'item_progress' && e.status === 'failed')
|
||||
expect(done).toHaveLength(2)
|
||||
expect(failed).toHaveLength(1)
|
||||
expect(events.at(-1)).toMatchObject({ type: 'batch_complete', succeeded: 2, failed: 1 })
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
export interface SseWorkerResult {
|
||||
ok: boolean
|
||||
warnings?: string[]
|
||||
error?: string
|
||||
}
|
||||
|
||||
export interface RunSseBatchOptions<T> {
|
||||
items: T[]
|
||||
concurrency: number
|
||||
/** Label used as the `classname` field on item_progress (kept for the existing client parser). */
|
||||
labelOf: (item: T) => string
|
||||
/** Per-item worker. `report(status)` streams intermediate progress (e.g. 'downloading'). */
|
||||
worker: (item: T, index: number, report: (status: string) => void) => Promise<SseWorkerResult>
|
||||
}
|
||||
|
||||
/**
|
||||
* Generic SSE batch runner. Emits the same event shape the furni client
|
||||
* parser consumes: batch_start / item_progress / batch_complete.
|
||||
*/
|
||||
export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
|
||||
const { items, labelOf, worker } = opts
|
||||
const concurrency = Math.min(Math.max(opts.concurrency || 3, 1), 5)
|
||||
const encoder = new TextEncoder()
|
||||
|
||||
const stream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const send = (data: unknown) => {
|
||||
try {
|
||||
controller.enqueue(encoder.encode(`data: ${JSON.stringify(data)}\n\n`))
|
||||
} catch {
|
||||
/* stream closed by client */
|
||||
}
|
||||
}
|
||||
|
||||
const startTime = Date.now()
|
||||
send({ type: 'batch_start', total: items.length, concurrency })
|
||||
|
||||
let succeeded = 0
|
||||
let failed = 0
|
||||
let withWarnings = 0
|
||||
|
||||
for (let i = 0; i < items.length; i += concurrency) {
|
||||
const chunk = items.slice(i, i + concurrency)
|
||||
await Promise.allSettled(
|
||||
chunk.map(async (item, chunkIdx) => {
|
||||
const index = i + chunkIdx
|
||||
const classname = labelOf(item)
|
||||
send({ type: 'item_progress', classname, status: 'started', index })
|
||||
try {
|
||||
const result = await worker(item, index, (status) =>
|
||||
send({ type: 'item_progress', classname, status, index }),
|
||||
)
|
||||
if (result.ok) {
|
||||
succeeded++
|
||||
if (result.warnings?.length) withWarnings++
|
||||
send({
|
||||
type: 'item_progress',
|
||||
classname,
|
||||
status: 'done',
|
||||
index,
|
||||
warnings: result.warnings?.length ? result.warnings : undefined,
|
||||
})
|
||||
} else {
|
||||
failed++
|
||||
send({
|
||||
type: 'item_progress',
|
||||
classname,
|
||||
status: 'failed',
|
||||
index,
|
||||
message: result.error,
|
||||
})
|
||||
}
|
||||
} catch (err) {
|
||||
failed++
|
||||
send({
|
||||
type: 'item_progress',
|
||||
classname,
|
||||
status: 'failed',
|
||||
index,
|
||||
message: (err as Error).message,
|
||||
})
|
||||
}
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
send({
|
||||
type: 'batch_complete',
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
duration: Date.now() - startTime,
|
||||
})
|
||||
controller.close()
|
||||
},
|
||||
})
|
||||
|
||||
return new Response(stream, {
|
||||
headers: {
|
||||
'Content-Type': 'text/event-stream',
|
||||
'Cache-Control': 'no-cache',
|
||||
Connection: 'keep-alive',
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user