1
/** Content-addressed, owner-private local attachment storage. */3
import { createHash, randomUUID } from 'node:crypto'4
import { constants, createReadStream } from 'node:fs'5
import { chmod, link, mkdir, open, readFile, unlink } from 'node:fs/promises'6
import { dirname, join, parse, resolve } from 'node:path'7
import {8
AttachmentError,9
AttachmentId,10
} from '@deepseek-ai/dsh-attachment'11
import type {12
ImageAttachmentLimits,13
ImageAttachmentRef,14
SaveImageAttachment,15
StoredImageAttachment,16
} from '@deepseek-ai/dsh-attachment'17
import { normalizeImage } from './normalization.ts'18
import type { NormalizationPolicy } from './normalization.ts'19
import { detectImage, probeImage } from './image.ts'20
import type { DetectedImage } from './image.ts'22
const ID_PATTERN = /^sha256:([a-f0-9]{64})$/23
const durableHomes = new Set<string>()25
function digest(data: Uint8Array): string {26
return createHash('sha256').update(data).digest('hex')27
}29
function displayName(value: string | undefined): string | undefined {30
if (value === undefined) return undefined31
// Strip both separator styles by hand: a POSIX host treats `\` as an32
// ordinary character, so path.basename would keep a Windows client's full33
// 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 : clean37
}39
function 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
}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
*/51
export function normalizedImagePath(root: string, ref: ImageAttachmentRef): string {52
const sha256 = ensureReference(ref)53
return join(root, 'objects', sha256.slice(0, 2), sha256)54
}56
async 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 detected65
}67
/**68
* Run the full admission policy for one image without touching storage,69
* including normalization: a batch whose members all validate cannot later70
* 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
*/76
export async function validateImageFile(77
input: SaveImageAttachment,78
limits: ImageAttachmentLimits,79
policy: NormalizationPolicy,80
): Promise<void> {81
await prepareImageFile(input, limits, policy)82
}84
/** Fully prepared normalized object, verified before any batch member is persisted. */85
export interface PreparedImageFile {86
/** Deterministic normalized bytes whose digest is {@link ref.attachmentId}. */87
data: Uint8Array88
/** Durable reference describing {@link data}. */89
ref: ImageAttachmentRef90
}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
*/99
export 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.height112
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
}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 never129
* reached storage, so the publication directory is synced before a durable130
* reference is reported.131
*/132
async 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') return135
/* 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
}145
/**146
* Create one private directory tree and persist every ancestor entry up to a147
* caller-vouched durable boundary. The walk deliberately ignores what mkdir148
* reports as newly created: a concurrent first save can create a level this149
* 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 a151
* directory the session checkpoint already references. Re-syncing a durable152
* 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
*/156
async 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 = target162
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) return167
level = parent168
}169
}171
/**172
* Establish this process's proof that one DSH_HOME entry and every ancestor173
* below the filesystem root are durable. Mere existence is insufficient: a174
* concurrent process may have created the directory but not synced its parent.175
*/176
async 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 home183
}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
*/191
export async function commitPreparedImageFile(192
root: string,193
prepared: PreparedImageFile,194
): Promise<ImageAttachmentRef> {195
const normalized = prepared.data196
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.ref202
}204
/**205
* Publish one immutable content-addressed object below a versioned attachment206
* root: staged write, fsync, hard-link into place, digest-verified EEXIST207
* deduplication, read-only mode, and durable directory entries from the208
* 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
*/214
export 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 data222
})())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
}230
/** Digest and byte count produced while streaming one immutable object to disk. */231
export interface StreamedImmutableObject {232
readonly sha256: string233
readonly bytes: number234
}236
/**237
* Stream one immutable object from bounded chunks into a staging file, then238
* 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
*/245
export 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: string253
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 error259
/* v8 ignore stop */260
}261
await publishStagedObject(root, target, staged)262
return { sha256: staged.sha256, bytes: staged.bytes }263
}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
*/272
export 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 error287
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) break297
}298
} catch (error) {299
if (error instanceof AttachmentError) throw error300
throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })301
}302
}304
interface StagedImmutableObject extends StreamedImmutableObject {305
readonly path: string306
readonly boundary: string307
}309
async 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 directory317
// 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 handle322
try {323
handle = await open(temporary, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600)324
const hash = createHash('sha256')325
let bytes = 0326
for await (const chunk of data) {327
signal?.throwIfAborted()328
await handle.writeFile(chunk)329
hash.update(chunk)330
bytes += chunk.byteLength331
}332
signal?.throwIfAborted()333
await handle.sync()334
signal?.throwIfAborted()335
await handle.close()336
handle = undefined337
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 error346
throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })347
}348
}350
async 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 error363
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 to368
// 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 restores371
// 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-creation374
// window before the reference can reach a session checkpoint. The dedup375
// path repeats these syncs because it may observe another writer's link376
// 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) break382
}383
} catch (error) {384
await removeTemporary(staged.path)385
if (error instanceof AttachmentError) throw error386
throw new AttachmentError('Unable to persist attachment.', 'ATTACHMENT_WRITE_FAILED', { cause: error })387
}388
}390
async 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
}396
async 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 cleanupError402
},403
)404
}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
*/414
export 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
}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
*/431
export 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: Uint8Array439
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, so449
// the read path only re-derives the header fields (no raster decode, no450
// 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.bytes454
|| 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
}