d3ro-voice/apps/desktop/src/main/services/runtime/download-part.ts

194 lines
7.5 KiB
TypeScript

// src/main/services/runtime/download-part.ts
// 런타임 부품 하나를 디스크로 내려받고 디스크 기준으로 크기·해시를 검증한다.
//
// 타임아웃은 "전송 전체" 가 아니라 "데이터가 멈춘 시간" 을 잰다. 예전에는
// AbortSignal.timeout(120s) 이 본문 전송 전체를 덮어서, 90MiB 부품을 120초 안에 받을 수
// 없는 링크(약 6Mbps 미만 — 테더링, 붐비는 Wi-Fi)에서는 세 번 모두 같은 자리에서 끊겨
// 로컬 엔진을 영영 설치할 수 없었다. 이제는 청크가 도착할 때마다 타이머가 다시 시작된다.
//
// 재시도 시 이미 받은 바이트가 있으면 Range 로 이어 받는다. 서버가 206 을 주지 않으면
// 처음부터 다시 받는다. 크기·해시가 틀리면 부분 파일을 지우고 처음부터 받는다.
import { createReadStream, createWriteStream } from 'node:fs'
import { rm, stat } from 'node:fs/promises'
import { createHash } from 'node:crypto'
import { Readable, Transform, Writable } from 'node:stream'
import type { ReadableStream as NodeReadableStream } from 'node:stream/web'
import { pipeline } from 'node:stream/promises'
import { D3ROError, ErrorCode } from '@d3ro/core/errors'
import type { RuntimePart } from './runtime-index'
/** fetch 응답 중 이 모듈이 쓰는 부분만 — 테스트에서 가짜 fetch 를 주입할 수 있게 한다 */
export interface RuntimeFetchResponse {
ok: boolean
status: number
headers: { get(name: string): string | null }
body: ReadableStream<Uint8Array> | null
json(): Promise<unknown>
}
export type RuntimeFetch = (
url: string,
init: { signal: AbortSignal; headers?: Record<string, string> },
) => Promise<RuntimeFetchResponse>
export const DEFAULT_STALL_TIMEOUT_MS = 60_000
/** 부품 다운로드 시도 횟수 — 전송 중 잘림/일시적 네트워크 오류 대비 */
export const DEFAULT_PART_DOWNLOAD_ATTEMPTS = 3
export interface DownloadPartOptions {
fetchImpl: RuntimeFetch
/** 이 시간 동안 새 데이터가 한 바이트도 오지 않으면 전송을 끊는다 (응답 헤더 대기 포함) */
stallTimeoutMs?: number
attempts?: number
/** 시도 하나가 실패했을 때 (로그용) */
onAttemptFailed?: (message: string, attempt: number, attempts: number) => void
}
/** 파일 SHA-256 (스트림 — 메모리에 통째로 올리지 않는다) */
export async function sha256File(path: string): Promise<string> {
const hash = createHash('sha256')
await pipeline(
createReadStream(path),
new Writable({
write(chunk: Buffer, _encoding: BufferEncoding, callback: (error?: Error | null) => void) {
hash.update(chunk)
callback()
},
}),
)
return hash.digest('hex')
}
async function fileSize(path: string): Promise<number> {
try {
return (await stat(path)).size
} catch {
return 0
}
}
/** `Content-Range: bytes <start>-<end>/<total>` 의 시작 오프셋 */
function contentRangeStart(header: string | null): number | null {
if (!header) return null
const match = /^bytes\s+(\d+)-\d+\/(?:\d+|\*)$/i.exec(header.trim())
return match ? Number(match[1]) : null
}
function failure(message: string): D3ROError {
return new D3ROError(ErrorCode.STTSidecarSpawnFailed, message)
}
/**
* 시도 한 번: 응답을 받아 파일에 쓴다 (이어 받기 가능하면 이어 쓴다).
* 데이터가 stallTimeoutMs 동안 멈추면 요청과 스트림을 모두 끊는다.
*/
async function transferOnce(
part: RuntimePart,
partPath: string,
fetchImpl: RuntimeFetch,
stallTimeoutMs: number,
): Promise<void> {
const offset = await fileSize(partPath)
const resumeFrom = offset > 0 && offset < part.size ? offset : 0
if (offset > 0 && resumeFrom === 0) {
await rm(partPath, { force: true })
}
const controller = new AbortController()
let source: Readable | null = null
let timer: NodeJS.Timeout | undefined
const stallError = failure(
`런타임 부품 전송이 ${Math.round(stallTimeoutMs / 1000)}초 동안 멈췄습니다 (${part.name})`,
)
const arm = (): void => {
if (timer) clearTimeout(timer)
timer = setTimeout(() => {
controller.abort(stallError)
source?.destroy(stallError)
}, stallTimeoutMs)
}
arm()
try {
const headers: Record<string, string> | undefined =
resumeFrom > 0 ? { Range: `bytes=${resumeFrom}-` } : undefined
const response = await fetchImpl(part.url, { signal: controller.signal, headers })
if (!response.ok || !response.body) {
if (resumeFrom > 0) await rm(partPath, { force: true })
throw failure(`런타임 부품을 받을 수 없습니다 (HTTP ${response.status}): ${part.name}`)
}
let append = false
if (resumeFrom > 0 && response.status === 206) {
if (contentRangeStart(response.headers.get('content-range')) !== resumeFrom) {
await rm(partPath, { force: true })
throw failure(`런타임 부품 이어 받기 범위가 맞지 않습니다 (${part.name})`)
}
append = true
}
source = Readable.fromWeb(response.body as NodeReadableStream<Uint8Array>)
const watchdog = new Transform({
transform(chunk: Buffer, _encoding: BufferEncoding, callback: (error?: Error | null, data?: Buffer) => void) {
arm()
callback(null, chunk)
},
})
// 스트림을 파일로 저장한 뒤 "디스크에 실제로 남은 파일"에서 크기와 해시를 계산한다.
// 메모리 스트림에서 센 값으로 검증하면, 디스크 쓰기가 잘려도 부품 검사를 통과해
// 결합 단계에 가서야 해시 불일치로 터진다 — 실측 사고.
await pipeline(source, watchdog, createWriteStream(partPath, { flags: append ? 'a' : 'w' }))
} catch (err) {
if (controller.signal.aborted) throw stallError
throw err
} finally {
if (timer) clearTimeout(timer)
}
}
/**
* 부품 하나를 partPath 로 내려받고 크기·해시를 확인한다. 받은 바이트 수를 돌려준다.
* 모든 시도가 실패하면 부분 파일을 지우고 마지막 오류를 던진다.
*/
export async function downloadPart(
part: RuntimePart,
partPath: string,
options: DownloadPartOptions,
): Promise<number> {
const stallTimeoutMs = options.stallTimeoutMs ?? DEFAULT_STALL_TIMEOUT_MS
const attempts = Math.max(1, options.attempts ?? DEFAULT_PART_DOWNLOAD_ATTEMPTS)
let lastError: Error | null = null
await rm(partPath, { force: true })
for (let attempt = 1; attempt <= attempts; attempt += 1) {
try {
await transferOnce(part, partPath, options.fetchImpl, stallTimeoutMs)
const partSize = await fileSize(partPath)
if (partSize < part.size) {
// 전송이 오류 없이 짧게 끝났다 — 다음 시도에서 이어 받는다
throw failure(`런타임 부품 크기 불일치 (${part.name}: ${partSize} != ${part.size})`)
}
if (partSize !== part.size) {
await rm(partPath, { force: true })
throw failure(`런타임 부품 크기 불일치 (${part.name}: ${partSize} != ${part.size})`)
}
const actualHash = await sha256File(partPath)
if (actualHash !== part.sha256) {
await rm(partPath, { force: true })
throw failure(`런타임 부품 해시 불일치 (${part.name})`)
}
return partSize
} catch (err) {
lastError = err instanceof Error ? err : new Error(String(err))
options.onAttemptFailed?.(lastError.message, attempt, attempts)
}
}
await rm(partPath, { force: true }).catch(() => undefined)
throw lastError ?? failure(`런타임 부품 다운로드 실패 (${part.name})`)
}