返回源码地图

packages/attachment/attachment-local/src/store.ts

main snapshot · da00f7f5358f · 正文引用章节 15 / 15;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Content-addressed, owner-private local attachment storage. */
2
3import { createHash, randomUUID } from 'node:crypto'
4import { constants, createReadStream } from 'node:fs'
5import { chmod, link, mkdir, open, readFile, unlink } from 'node:fs/promises'
6import { dirname, join, parse, resolve } from 'node:path'
7import {
8 AttachmentError,
9 AttachmentId,
10} from '@deepseek-ai/dsh-attachment'
11import type {
12 ImageAttachmentLimits,
13 ImageAttachmentRef,
14 SaveImageAttachment,
15 StoredImageAttachment,
16} from '@deepseek-ai/dsh-attachment'
17import { normalizeImage } from './normalization.ts'
18import type { NormalizationPolicy } from './normalization.ts'
19import { detectImage, probeImage } from './image.ts'
20import type { DetectedImage } from './image.ts'
21
22const ID_PATTERN = /^sha256:([a-f0-9]{64})$/
23const durableHomes = new Set<string>()
24
25function digest(data: Uint8Array): string {
26 return createHash('sha256').update(data).digest('hex')
27}
28
29function displayName(value: string | undefined): string | undefined {
30 if (value === undefined) return undefined
31 // Strip both separator styles by hand: a POSIX host treats `\` as an
32 // ordinary character, so path.basename would keep a Windows client's full
33 // local path and leak it into the reference and the session log.
34 const leaf = value.slice(Math.max(value.lastIndexOf('/'), value.lastIndexOf('\\')) + 1)
35 const clean = leaf.replace(/[\u0000-\u001f\u007f]/g, '').trim().slice(0, 255)
36 return clean === '' ? undefined : clean
37}
38
39function ensureReference(ref: ImageAttachmentRef): string {
40 const match = ID_PATTERN.exec(String(ref.attachmentId))
41 if (match?.[1] === undefined) throw new AttachmentError('Attachment reference is invalid.', 'INVALID_ATTACHMENT_REF')
42 return match[1]
43}
44
45/**
46 * Derive the absolute immutable-object path for one normalized attachment.
47 * @param root - absolute `DSH_HOME/attachments/v1` root.
48 * @param ref - durable normalized attachment reference.
49 * @returns provider-local path without reading the object.
50 */
51export function normalizedImagePath(root: string, ref: ImageAttachmentRef): string {
52 const sha256 = ensureReference(ref)
53 return join(root, 'objects', sha256.slice(0, 2), sha256)
54}
55
56async function inspectMetadata(
57 data: Uint8Array,
58 declaredMediaType: ImageAttachmentRef['mediaType'],
59 limits: ImageAttachmentLimits,
60): Promise<DetectedImage> {
61 if (data.byteLength === 0) throw new AttachmentError('Image is empty.', 'INVALID_IMAGE')
62 const detected = await detectImage(data, { maxPixels: limits.maxImagePixels, maxDimension: limits.maxImageDimension })
63 if (detected.mediaType !== declaredMediaType) throw new AttachmentError('Declared image type does not match its bytes.', 'IMAGE_TYPE_MISMATCH')
64 return detected
65}
66
67/**
68 * Run the full admission policy for one image without touching storage,
69 * including normalization: a batch whose members all validate cannot later
70 * be refused by the normalized image byte cap during publication.
71 * @param input - encoded bytes and declared metadata.
72 * @param limits - resolved source admission policy.
73 * @param policy - resolved normalization policy.
74 * @returns completion after the raster has been decoded and its normalized version proven to fit.
75 */
76export async function validateImageFile(
77 input: SaveImageAttachment,
78 limits: ImageAttachmentLimits,
79 policy: NormalizationPolicy,
80): Promise<void> {
81 await prepareImageFile(input, limits, policy)
82}
83
84/** Fully prepared normalized object, verified before any batch member is persisted. */
85export interface PreparedImageFile {
86 /** Deterministic normalized bytes whose digest is {@link ref.attachmentId}. */
87 data: Uint8Array
88 /** Durable reference describing {@link data}. */
89 ref: ImageAttachmentRef
90}
91
92/**
93 * Decode, normalize, and verify one submitted image without touching storage.
94 * @param input - submitted encoded bytes and declared media type.
95 * @param limits - source admission policy.
96 * @param policy - independent normalization policy.
97 * @returns immutable reference facts beside bytes ready for atomic publication.
98 */
99export async function prepareImageFile(
100 input: SaveImageAttachment,
101 limits: ImageAttachmentLimits,
102 policy: NormalizationPolicy,
103): Promise<PreparedImageFile> {
104 if (input.data.byteLength > limits.maxImageBytes) {
105 throw new AttachmentError('Image exceeds the configured byte limit.', 'IMAGE_TOO_LARGE')
106 }
107 const detected = await inspectMetadata(input.data, input.mediaType, limits)
108 const normalized = await normalizeImage(input.data, detected, policy)
109 const sha256 = digest(normalized.data)
110 const name = displayName(input.name)
111 const downscaled = detected.width !== normalized.width || detected.height !== normalized.height
112 return {
113 data: normalized.data,
114 ref: {
115 attachmentId: AttachmentId(`sha256:${sha256}`),
116 mediaType: normalized.mediaType,
117 width: normalized.width,
118 height: normalized.height,
119 bytes: normalized.data.byteLength,
120 ...(name !== undefined ? { name } : {}),
121 ...downscaled ? { originalDimensions: { width: detected.width, height: detected.height } } : {},
122 },
123 }
124}
125
126/**
127 * Make a directory's entries durable (fsync on a read-only directory handle).
128 * A synced file alone does not survive a crash when its directory entry never
129 * reached storage, so the publication directory is synced before a durable
130 * reference is reported.
131 */
132async function syncDirectory(path: string): Promise<void> {
133 /* v8 ignore next -- Windows cannot open directory handles; NTFS metadata journaling owns entry durability there. */
134 if (process.platform === 'win32') return
135 /* v8 ignore start -- Windows cannot exercise directory fsync; POSIX behavior tests enforce this peer. */
136 const handle = await open(path, constants.O_RDONLY)
137 try {
138 await handle.sync()
139 } finally {
140 await handle.close()
141 }
142 /* v8 ignore stop */
143}
144
145/**
146 * Create one private directory tree and persist every ancestor entry up to a
147 * caller-vouched durable boundary. The walk deliberately ignores what mkdir
148 * reports as newly created: a concurrent first save can create a level this
149 * process then merely observes, so "already existed" is not "already durable"
150 * — the entry may still be unsynced in the creator, and a crash would drop a
151 * directory the session checkpoint already references. Re-syncing a durable
152 * entry is harmless; skipping an unsynced one is not.
153 * @param path - absolute directory to create.
154 * @param boundary - absolute ancestor the caller vouches is already durable.
155 */
156async function ensureDurableDirectory(path: string, boundary: string): Promise<void> {
157 const target = resolve(path)
158 const stop = resolve(boundary)
159 await mkdir(target, { recursive: true, mode: 0o700 })
160 await chmod(target, 0o700)
161 let level = target
162 while (level !== stop) {
163 const parent = dirname(level)
164 await syncDirectory(parent)
165 /* v8 ignore next -- filesystem-root guard: callers pass a boundary that is an ancestor of path, so the walk reaches it first. */
166 if (parent === level) return
167 level = parent
168 }
169}
170
171/**
172 * Establish this process's proof that one DSH_HOME entry and every ancestor
173 * below the filesystem root are durable. Mere existence is insufficient: a
174 * concurrent process may have created the directory but not synced its parent.
175 */
176async function ensureDurableHome(path: string): Promise<string> {
177 const home = resolve(path)
178 if (!durableHomes.has(home)) {
179 await ensureDurableDirectory(home, parse(home).root)
180 durableHomes.add(home)
181 }
182 return home
183}
184
185/**
186 * Publish one already verified normalized image below a versioned attachment root.
187 * @param root - absolute `DSH_HOME/attachments/v1` root.
188 * @param prepared - deterministic normalized bytes and reference.
189 * @returns durable content-addressed normalized image reference.
190 */
191export async function commitPreparedImageFile(
192 root: string,
193 prepared: PreparedImageFile,
194): Promise<ImageAttachmentRef> {
195 const normalized = prepared.data
196 const sha256 = ensureReference(prepared.ref)
197 if (digest(normalized) !== sha256 || normalized.byteLength !== prepared.ref.bytes) {
198 throw new AttachmentError('Prepared attachment bytes do not match their reference.', 'ATTACHMENT_CORRUPT')
199 }
200 await publishImmutableObject(root, normalizedImagePath(root, prepared.ref), normalized, sha256)
201 return prepared.ref
202}
203
204/**
205 * Publish one immutable content-addressed object below a versioned attachment
206 * root: staged write, fsync, hard-link into place, digest-verified EEXIST
207 * deduplication, read-only mode, and durable directory entries from the
208 * target's parent up to (excluding) `root`.
209 * @param root - absolute `DSH_HOME/attachments/v1` root.
210 * @param target - absolute final object path below `root`.
211 * @param data - exact object bytes whose digest is `sha256`.
212 * @param sha256 - hex digest the stored bytes must match on deduplication.
213 */
214export async function publishImmutableObject(
215 root: string,
216 target: string,
217 data: Uint8Array,
218 sha256: string,
219): Promise<void> {
220 const staged = await stageImmutableObject(root, (function* (): Iterable<Uint8Array> {
221 yield data
222 })())
223 if (staged.sha256 !== sha256) {
224 await removeTemporary(staged.path)
225 throw new AttachmentError('Attachment bytes do not match their publication digest.', 'ATTACHMENT_CORRUPT')
226 }
227 await publishStagedObject(root, target, staged)
228}
229
230/** Digest and byte count produced while streaming one immutable object to disk. */
231export interface StreamedImmutableObject {
232 readonly sha256: string
233 readonly bytes: number
234}
235
236/**
237 * Stream one immutable object from bounded chunks into a staging file, then
238 * publish it at a digest-derived target without collecting the complete object in memory.
239 * @param root - absolute `DSH_HOME/attachments/v1` root.
240 * @param data - exact object bytes in order.
241 * @param targetFor - derive the final absolute target from the completed digest and byte count.
242 * @param signal - optional cancellation for source reads and storage writes.
243 * @returns digest and exact byte count of the published object.
244 */
245export async function publishImmutableObjectStream(
246 root: string,
247 data: AsyncIterable<Uint8Array>,
248 targetFor: (sha256: string, bytes: number) => string,
249 signal?: AbortSignal,
250): Promise<StreamedImmutableObject> {
251 const staged = await stageImmutableObject(root, data, signal)
252 let target: string
253 try {
254 target = targetFor(staged.sha256, staged.bytes)
255 } catch (error) {
256 /* v8 ignore start -- The local target callback constructs a validated reference from this function's digest. */
257 await removeTemporary(staged.path)
258 throw error
259 /* v8 ignore stop */
260 }
261 await publishStagedObject(root, target, staged)
262 return { sha256: staged.sha256, bytes: staged.bytes }
263}
264
265/**
266 * Publish another durable hard-link name for an existing immutable object.
267 * @param root - absolute versioned attachment root.
268 * @param source - existing content-addressed object below `root`.
269 * @param target - new alias below `root`.
270 * @param sha256 - expected object digest for an existing-target race.
271 */
272export async function publishImmutableAlias(
273 root: string,
274 source: string,
275 target: string,
276 sha256: string,
277): Promise<void> {
278 const parent = dirname(target)
279 try {
280 const boundary = await ensureDurableHome(dirname(dirname(resolve(root))))
281 await ensureDurableDirectory(parent, boundary)
282 try {
283 await link(source, target)
284 } catch (error) {
285 /* v8 ignore next -- Private same-filesystem directories make EEXIST the only recoverable link race. */
286 if (!(error instanceof Error && 'code' in error && error.code === 'EEXIST')) throw error
287 if (await digestFile(target) !== sha256) {
288 throw new AttachmentError('Stored attachment failed integrity verification.', 'ATTACHMENT_CORRUPT')
289 }
290 }
291 await chmod(target, 0o400)
292 const stop = resolve(root)
293 for (let level = parent; level !== stop; level = dirname(level)) {
294 await syncDirectory(level)
295 /* v8 ignore next -- filesystem-root guard: targets sit below root, so the walk reaches `stop` first. */
296 if (dirname(level) === level) break
297 }
298 } catch (error) {
299 if (error instanceof AttachmentError) throw error
300 throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })
301 }
302}
303
304interface StagedImmutableObject extends StreamedImmutableObject {
305 readonly path: string
306 readonly boundary: string
307}
308
309async function stageImmutableObject(
310 root: string,
311 data: AsyncIterable<Uint8Array> | Iterable<Uint8Array>,
312 signal?: AbortSignal,
313): Promise<StagedImmutableObject> {
314 const staging = join(root, 'tmp')
315 // Establish DSH_HOME itself against the filesystem root once per process.
316 // Every process performs that proof independently, so observing a directory
317 // another process created can never be mistaken for durable publication.
318 const boundary = await ensureDurableHome(dirname(dirname(resolve(root))))
319 await ensureDurableDirectory(staging, boundary)
320 const temporary = join(staging, randomUUID())
321 let handle
322 try {
323 handle = await open(temporary, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600)
324 const hash = createHash('sha256')
325 let bytes = 0
326 for await (const chunk of data) {
327 signal?.throwIfAborted()
328 await handle.writeFile(chunk)
329 hash.update(chunk)
330 bytes += chunk.byteLength
331 }
332 signal?.throwIfAborted()
333 await handle.sync()
334 signal?.throwIfAborted()
335 await handle.close()
336 handle = undefined
337 return { path: temporary, boundary, sha256: hash.digest('hex'), bytes }
338 } catch (error) {
339 /* v8 ignore next -- A descriptor remains open only when write, sync, or close fails. */
340 if (handle !== undefined) await handle.close().catch(
341 /* v8 ignore next -- Close failure is superseded by the storage operation that entered cleanup. */
342 () => {},
343 )
344 await removeTemporary(temporary)
345 if (error instanceof AttachmentError || signal?.aborted === true) throw error
346 throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })
347 }
348}
349
350async function publishStagedObject(
351 root: string,
352 target: string,
353 staged: StagedImmutableObject,
354): Promise<void> {
355 const parent = dirname(target)
356 try {
357 await ensureDurableDirectory(parent, staged.boundary)
358 try {
359 await link(staged.path, target)
360 } catch (error) {
361 /* v8 ignore next -- Private same-filesystem directories make EEXIST the only recoverable link race. */
362 if (!(error instanceof Error && 'code' in error && error.code === 'EEXIST')) throw error
363 if (await digestFile(target) !== staged.sha256) {
364 throw new AttachmentError('Stored attachment failed integrity verification.', 'ATTACHMENT_CORRUPT')
365 }
366 }
367 // Windows shares the read-only attribute across hard links and refuses to
368 // unlink either name once it is set, so discard the staging name first.
369 await unlink(staged.path)
370 // The target remains the sole link for a new object; this also restores
371 // read-only mode when the deduplication path observes an existing object.
372 await chmod(target, 0o400)
373 // Persist the target entry and close every concurrent parent-creation
374 // window before the reference can reach a session checkpoint. The dedup
375 // path repeats these syncs because it may observe another writer's link
376 // before that writer reaches its own durability boundary.
377 const stop = resolve(root)
378 for (let level = parent; level !== stop; level = dirname(level)) {
379 await syncDirectory(level)
380 /* v8 ignore next -- filesystem-root guard: targets sit below root, so the walk reaches `stop` first. */
381 if (dirname(level) === level) break
382 }
383 } catch (error) {
384 await removeTemporary(staged.path)
385 if (error instanceof AttachmentError) throw error
386 throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })
387 }
388}
389
390async function digestFile(path: string): Promise<string> {
391 const hash = createHash('sha256')
392 for await (const chunk of createReadStream(path) as AsyncIterable<Buffer>) hash.update(chunk)
393 return hash.digest('hex')
394}
395
396async function removeTemporary(path: string): Promise<void> {
397 await unlink(path).catch(
398 /* v8 ignore next -- Cleanup can observe a staging name already removed after successful linking. */
399 (cleanupError: unknown) => {
400 /* v8 ignore next -- Any cleanup failure except an absent staging name must remain visible. */
401 if (!(cleanupError instanceof Error && 'code' in cleanupError && cleanupError.code === 'ENOENT')) throw cleanupError
402 },
403 )
404}
405
406/**
407 * Decode and normalize one image once, then publish the prepared object.
408 * @param root - absolute `DSH_HOME/attachments/v1` root.
409 * @param input - submitted encoded bytes and declared media type.
410 * @param limits - resolved source admission policy.
411 * @param policy - resolved normalization policy.
412 * @returns durable content-addressed normalized image reference.
413 */
414export async function saveImageFile(
415 root: string,
416 input: SaveImageAttachment,
417 limits: ImageAttachmentLimits,
418 policy: NormalizationPolicy,
419): Promise<ImageAttachmentRef> {
420 return commitPreparedImageFile(root, await prepareImageFile(input, limits, policy))
421}
422
423/**
424 * Read and verify one content-addressed image.
425 * @param root - absolute `DSH_HOME/attachments/v1` root.
426 * @param ref - reference recorded in the session log.
427 * @param signal - optional cancellation for filesystem and verification work.
428 * @returns verified bytes and reference.
429 * @throws the signal reason when aborted, or an AttachmentError when verification fails.
430 */
431export async function readImageFile(
432 root: string,
433 ref: ImageAttachmentRef,
434 signal?: AbortSignal,
435): Promise<StoredImageAttachment> {
436 signal?.throwIfAborted()
437 const sha256 = ensureReference(ref)
438 let data: Uint8Array
439 try {
440 data = new Uint8Array(await readFile(normalizedImagePath(root, ref), { signal }))
441 } catch (error) {
442 signal?.throwIfAborted()
443 if (error instanceof Error && 'code' in error && error.code === 'ENOENT') throw new AttachmentError('Attachment object is missing.', 'ATTACHMENT_NOT_FOUND')
444 throw new AttachmentError('Unable to read image attachment.', 'ATTACHMENT_READ_FAILED', { cause: error })
445 }
446 signal?.throwIfAborted()
447 if (digest(data) !== sha256) throw new AttachmentError('Stored attachment failed integrity verification.', 'ATTACHMENT_CORRUPT')
448 // The digest proves these are the exact bytes admission fully decoded, so
449 // the read path only re-derives the header fields (no raster decode, no
450 // per-request pixel amplification on history replay).
451 const metadata = await probeImage(data)
452 signal?.throwIfAborted()
453 if (metadata.mediaType !== ref.mediaType || data.byteLength !== ref.bytes
454 || metadata.width !== ref.width || metadata.height !== ref.height) {
455 throw new AttachmentError('Stored attachment metadata does not match its reference.', 'ATTACHMENT_CORRUPT')
456 }
457 return { ref, data }
458}