diff --git a/llmdoc/meta.json b/llmdoc/meta.json index 9f1d8f6a..512c65ae 100644 --- a/llmdoc/meta.json +++ b/llmdoc/meta.json @@ -66,7 +66,7 @@ "validatedRevision": "2b45dc99c607f59351d214845fe8e0c6b912f478" }, "store/storage-backends.mdx": { - "validatedRevision": "2b45dc99c607f59351d214845fe8e0c6b912f478" + "validatedRevision": "69eeb627035f9dc6501a316638a137bec986a8d5" } }, "convergence": { diff --git a/llmdoc/store/storage-backends.mdx b/llmdoc/store/storage-backends.mdx index b079c37b..94703bf6 100644 --- a/llmdoc/store/storage-backends.mdx +++ b/llmdoc/store/storage-backends.mdx @@ -1,5 +1,5 @@ --- -description: S3 后端管理:不可变 backendId、active 默认指针、Context 固定绑定、credential generation、真实能力探测与 socket DNS 安全。 +description: S3 后端管理:不可变 backendId、active 默认指针、Context 固定绑定、credential generation、未知长度上传的临时磁盘与 deadline、真实能力探测及 socket DNS 安全。 kind: architecture relations: related: @@ -29,6 +29,8 @@ code: - packages/dashboard/src/pages/system/forms/StorageConnectionFields.tsx - packages/dashboard/src/pages/system/forms/storageConnection.ts - packages/server/test/s3ObjectStore.integration.test.ts + - packages/server/test/s3Objects.test.ts + - packages/server/test/s3Probe.test.ts - packages/server/test/managedSecurity.test.ts --- @@ -61,6 +63,14 @@ Node 使用官方 AWS S3 SDK,配置 path-style、禁自动 redirect、单次 test 在随机隔离 namespace 执行真实 PUT/HEAD/GET/DELETE、metadata、空/流式对象、特殊字符 key、分页 continuation,并以并发 `If-None-Match` / `If-Match` 写证明对象服务的原子条件语义。错误 ETag 必須拒绝,probe cleanup 是激活条件之一。`HEAD → PUT` 不能模拟原子 compare-and-set;HTTP mock 或能创建 bucket 不构成兼容证据。 +分页检查使用独立子 namespace 与成功写入的预期集合,不能把其他检查的失败写入尝试当作分页缺失;cleanup 则仍清理所有尝试过的 key,包括写入后才报错的对象。 + +能力检查失败时,服务端 console 只记录固定检查名与严格白名单归一的原因;任意异常消息、cause、key、endpoint 和凭证均不能回显。返回的布尔检查结果仍决定激活资格,诊断用于区分检查语义失败与传输异常,不替代真实兼容验证。 + 标准 Node driver 当前不提供 presign,因此 default Store 使用流式 relay,平台 S3 Context 不广告可选 direct-upload。core/neutral Store client 仍定义严格 exact-size direct 契约供显式能力 driver 使用,但未实现的浏览器 CORS/直传不能标记已交付。当前是 single PUT,无 multipart/resumable。 +未知长度的上传流先以背压写入系统临时目录,完整接收后按真实 `Content-Length` 执行 single PUT,保留原子条件头。这样避免把整份对象缓存在进程堆上;byte cap 在接收过程中执行,同一个 deadline 覆盖暂存、PUT 与确认 HEAD。临时目录与文件权限分别为 `0700`、`0600`,正常完成或错误退出均清理;进程崩溃后的磁盘回收仍需部署环境承担。 + +这条路径需要可写且容量足够的临时磁盘,并增加完整接收后才能开始远端上传的等待。临时文件是传输暂存,不是部署级本地 ObjectStore;持久对象仍由绑定的 S3 后端保存。兼容判断继续以当前部署的真实全项探测为准,不能从单项布尔失败推定厂商能力或唯一根因。 + API、`tb storage` 与 Dashboard 存储后端页提供同权 list/get/add/test/activate/credential update/delete。Context CLI 的 backend 选择与 Dashboard 字段必须落到相同注册 payload;普通 SK 不因拥有 registry read/write 就获得后端管理权限。对字节资源的真实验证每轮有界,保留脱敏证据;不得把 probe namespace 或签名 URL 写入日志。 diff --git a/packages/server/package.json b/packages/server/package.json index 286c45b3..a697d29f 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/server", - "version": "0.22.0", + "version": "0.22.1", "description": "Self-hosted Tool Bridge Node server with PostgreSQL, S3 object storage and WebSocket device channels", "type": "module", "license": "MIT", diff --git a/packages/server/src/s3Objects.ts b/packages/server/src/s3Objects.ts index 18f25814..e70ef265 100644 --- a/packages/server/src/s3Objects.ts +++ b/packages/server/src/s3Objects.ts @@ -15,9 +15,14 @@ import { type ObjectStore, TBError, } from '@tool-bridge/core' +import { createReadStream, createWriteStream } from 'node:fs' import { NodeHttpHandler } from '@smithy/node-http-handler' +import { finished, pipeline } from 'node:stream/promises' /** Official AWS SDK protocol adapter, deliberately confined to the Node host. */ import { Readable, Transform } from 'node:stream' +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import { s3Network, type S3NetworkOptions } from './s3Network' export interface S3ObjectStoreOptions extends S3NetworkOptions { @@ -94,10 +99,11 @@ function bounded(source: Readable, maxBytes: number): Readable { return output } -function uploadBody( +async function uploadBody( body: ObjectBody, maxBytes: number, -): { body: Buffer | Readable, length?: number } { + signal: AbortSignal, +): Promise<{ body: Buffer | Readable, dispose(): Promise, length: number }> { if ( typeof body === 'string' || body instanceof Uint8Array @@ -110,16 +116,39 @@ function uploadBody( ? Buffer.from(body) : Buffer.from(body.buffer, body.byteOffset, body.byteLength) if (value.length > maxBytes) throw limitError() - return { body: value, length: value.length } + return { body: value, length: value.length, dispose: async () => {} } } const reader = body.getReader() + let total = 0 + let complete = false + let cancelRequested = false + const cancel = () => { + if (complete || cancelRequested) return + cancelRequested = true + void reader.cancel?.(signal.reason).catch(() => {}) + } + async function readNext() { + signal.throwIfAborted() + let onAbort: () => void = () => {} + const aborted = new Promise((_, reject) => { + onAbort = () => { + cancel() + reject(signal.reason) + } + signal.addEventListener('abort', onAbort, { once: true }) + }) + try { + return await Promise.race([reader.read(), aborted]) + } finally { + // A shared, never-resolved abort promise would retain a reaction per chunk. + signal.removeEventListener('abort', onAbort) + } + } const stream = Readable.from( (async function* () { - let complete = false - let total = 0 try { for (;;) { - const { done, value } = await reader.read() + const { done, value } = await readNext() if (done) { complete = true break @@ -130,17 +159,37 @@ function uploadBody( yield value } } finally { - if (!complete) await reader.cancel?.().catch(() => {}) - reader.releaseLock() + cancel() } })(), { objectMode: false, highWaterMark: 64 * 1024 }, ) - // Ensure cancellation reaches a reader whose read() is still pending. - stream.on('close', () => { - void reader.cancel?.().catch(() => {}) - }) - return { body: stream } + let directory: string | undefined + try { + // Use an exact Content-Length instead of relying on unknown-length chunked + // PUT support. Spool with backpressure rather than buffering maxBytes in RAM. + directory = await mkdtemp(join(tmpdir(), 'tb-s3-upload-')) + const path = join(directory, 'body') + await pipeline(stream, createWriteStream(path, { flags: 'wx', mode: 0o600 }), { signal }) + const file = createReadStream(path) + const uploadDirectory = directory + return { + body: file, + length: total, + async dispose() { + file.destroy() + await finished(file, { cleanup: true }).catch(() => {}) + await rm(uploadDirectory, { recursive: true, force: true }) + }, + } + } catch (error) { + stream.destroy() + if (directory) await rm(directory, { recursive: true, force: true }) + throw error + } finally { + cancel() + reader.releaseLock() + } } export function createS3ObjectStore( @@ -237,20 +286,21 @@ export function createS3ObjectStore( } } - const head = async (key: string): Promise => { + const head = async (key: string, signal?: AbortSignal): Promise => { try { - return await run('HEAD', async signal => + const execute = async (signal: AbortSignal) => meta( key, await client.send( new HeadObjectCommand({ Bucket: config.bucket, Key: key }), { abortSignal: signal }, ), - ), - ) + ) + return signal ? await execute(signal) : await run('HEAD', execute) } catch (error) { - if (isTBError(error) && error.code === 'not_found') return null - throw error + const normalized = s3Error('HEAD', error) + if (normalized.code === 'not_found') return null + throw normalized } } @@ -298,51 +348,45 @@ export function createS3ObjectStore( 'invalid_argument', 'S3 conditional write modes are mutually exclusive', ) - const upload = uploadBody(body, maxBytes) - let bodyError: Error | undefined - const uploadController = new AbortController() - if (upload.body instanceof Readable) - upload.body.on('error', (error) => { - bodyError = error - uploadController.abort() - }) - try { - await run('PUT', signal => - client.send( - new PutObjectCommand({ - Bucket: config.bucket, - Key: key, - Body: upload.body, - ContentLength: upload.length, - ContentType: opts?.contentType, - Metadata: opts?.metadata, - IfNoneMatch: opts?.ifNoneMatch, - IfMatch: + return run('PUT', async (signal) => { + const upload = await uploadBody(body, maxBytes, signal) + const uploadController = new AbortController() + if (upload.body instanceof Readable) + upload.body.on('error', () => uploadController.abort()) + try { + try { + await client.send( + new PutObjectCommand({ + Bucket: config.bucket, + Key: key, + Body: upload.body, + ContentLength: upload.length, + ContentType: opts?.contentType, + Metadata: opts?.metadata, + IfNoneMatch: opts?.ifNoneMatch, + IfMatch: opts?.ifMatchEtag !== undefined ? `"${opts.ifMatchEtag}"` : undefined, - }), - { abortSignal: AbortSignal.any([signal, uploadController.signal]) }, - ), - ) - } catch (error) { - if (bodyError && isTBError(bodyError)) throw bodyError - if ( - opts?.ifMatchEtag !== undefined - && isTBError(error) - && error.code === 'not_found' - ) - throw new TBError('conflict', 'S3 PUT condition failed') - throw error - } finally { - if (upload.body instanceof Readable) upload.body.destroy() - } - const stored = await head(key) - if (!stored) - throw new TBError('unavailable', 'S3 PUT was not observable by HEAD', { - retryable: true, - }) - return stored + }), + { abortSignal: AbortSignal.any([signal, uploadController.signal]) }, + ) + } catch (error) { + const normalized = s3Error('PUT', error) + if (opts?.ifMatchEtag !== undefined && normalized.code === 'not_found') + throw new TBError('conflict', 'S3 PUT condition failed') + throw normalized + } + const stored = await head(key, signal) + if (!stored) + throw new TBError('unavailable', 'S3 PUT was not observable by HEAD', { + retryable: true, + }) + return stored + } finally { + await upload.dispose() + } + }) }, async delete(key) { try { diff --git a/packages/server/src/s3Probe.ts b/packages/server/src/s3Probe.ts index 0f78b882..95972457 100644 --- a/packages/server/src/s3Probe.ts +++ b/packages/server/src/s3Probe.ts @@ -8,6 +8,15 @@ export interface S3ProbeResult { cleanupSucceeded: boolean } +function diagnosticReason(error: unknown): string { + if (!isTBError(error)) return 'unexpected_error' + // TBError can also originate from an upload source. Only known adapter messages + // are safe; never log arbitrary messages, causes, keys, endpoints or credentials. + const normalized = /^(?:S3 (?:PUT|HEAD|GET|LIST|DELETE) (?:failed(?: \([1-5]\d{2}\))?|was denied|object not found|condition failed)|S3 GET deadline exceeded|S3 PUT was not observable by HEAD)$/.exec(error.message) + // JS $ also matches before a final newline, so require a complete match. + return normalized?.[0] === error.message ? error.message : 'operation_failed' +} + /** Destructive only inside a fresh, unguessable probe namespace in the selected bucket. */ export async function probeS3ObjectStore( config: S3StoreConfig, @@ -40,11 +49,15 @@ export async function probeS3ObjectStore( name: string, run: () => Promise, ): Promise { + let reason = 'unexpected_result' try { checks[name] = await run() - } catch { + } catch (error) { checks[name] = false + reason = diagnosticReason(error) } + if (!checks[name]) + console.warn('S3 capability probe check failed', { check: name, reason }) } async function read(name: string): Promise { const object = await store.get(name) @@ -161,11 +174,20 @@ export async function probeS3ObjectStore( return true }) await check('pagination', async () => { + // Other checks may fail before or after their PUT reaches the backend. + // Keep listing expectations independent from the cleanup attempt ledger. + const paginationPrefix = `${prefix}pagination/` + const expected = new Set() + for (let index = 0; index < 5; index++) { + const target = key(`pagination/item-${index}`) + await store.put(target, `page-${index}`, { ifNoneMatch: '*' }) + expected.add(target) + } const found = new Set() const cursors = new Set() let cursor: string | undefined do { - const page = await store.list(prefix, { cursor, limit: 2 }) + const page = await store.list(paginationPrefix, { cursor, limit: 2 }) for (const item of page.items) { if (!('key' in item) || found.has(item.key)) return false found.add(item.key) @@ -177,8 +199,8 @@ export async function probeS3ObjectStore( } } while (cursor) return ( - found.size === keys.size - && [...keys].every(item => found.has(item)) + found.size === expected.size + && [...expected].every(item => found.has(item)) && cursors.size > 0 ) }) diff --git a/packages/server/test/s3Objects.test.ts b/packages/server/test/s3Objects.test.ts index ae54aeeb..23b064d2 100644 --- a/packages/server/test/s3Objects.test.ts +++ b/packages/server/test/s3Objects.test.ts @@ -3,9 +3,12 @@ import { type IncomingMessage, type ServerResponse, } from 'node:http' -import { afterEach, describe, expect, it } from 'vitest' -import { readStreamBytes } from '@tool-bridge/core' +import { bytesToObjectStream, readStreamBytes } from '@tool-bridge/core' +import { mkdtemp, readdir, rm, stat } from 'node:fs/promises' +import { afterEach, describe, expect, it, vi } from 'vitest' import { once } from 'node:events' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import { createS3ObjectStore, type S3ObjectStoreOptions, @@ -15,8 +18,16 @@ import { assertS3Address } from '../src/s3Network' const cleanup: Array<() => Promise | void> = [] afterEach(async () => { for (const close of cleanup.splice(0).reverse()) await close() + vi.unstubAllEnvs() }) +async function uploadTempDirectory() { + const directory = await mkdtemp(join(tmpdir(), 'tb-test-upload-')) + vi.stubEnv(process.platform === 'win32' ? 'TEMP' : 'TMPDIR', directory) + cleanup.push(() => rm(directory, { recursive: true, force: true })) + return directory +} + async function fixture( handler: (request: IncomingMessage, response: ServerResponse) => void, options: S3ObjectStoreOptions = {}, @@ -178,9 +189,12 @@ describe('Node S3 wire and network policy', () => { }) it('bounds streamed uploads and cancels the source', async () => { + const directory = await uploadTempDirectory() let cancelled = false + let calls = 0 const { store } = await fixture( (request, response) => { + calls++ request.resume() request.on('end', () => response.end()) }, @@ -198,6 +212,188 @@ describe('Node S3 wire and network policy', () => { }), ).rejects.toMatchObject({ code: 'invalid_argument' }) expect(cancelled).toBe(true) + expect(calls).toBe(0) + expect(await readdir(directory)).toEqual([]) + }) + + it('sends unknown-length streams as an exact-length single PUT with atomic conditions', async () => { + const directory = await uploadTempDirectory() + const payload = Buffer.from('多块 streamed content '.repeat(100)) + const requests: string[] = [] + let received = Buffer.alloc(0) + let headers: IncomingMessage['headers'] = {} + let duringUpload: Promise | undefined + const { store } = await fixture((request, response) => { + requests.push(request.method!) + if (request.method === 'HEAD') { + response.writeHead(200, { 'content-length': payload.length, 'etag': '"stored"' }) + response.end() + return + } + headers = request.headers + duringUpload = (async () => { + const [name] = await readdir(directory) + expect(name).toMatch(/^tb-s3-upload-/) + if (!name) throw new Error('upload did not spool to a private directory') + const path = join(directory, name) + if (process.platform !== 'win32') { + expect((await stat(path)).mode & 0o777).toBe(0o700) + expect((await stat(join(path, 'body'))).mode & 0o777).toBe(0o600) + } + expect((await stat(join(path, 'body'))).size).toBe(payload.length) + })() + const chunks: Buffer[] = [] + request.on('data', chunk => chunks.push(chunk)) + request.on('end', () => { + received = Buffer.concat(chunks) + void duringUpload!.then(() => response.end(), () => { + response.writeHead(500) + response.end() + }) + }) + }) + let offset = 0 + const result = await store.put('stream', { + getReader: () => ({ + read: async () => { + if (offset === payload.length) return { done: true } + const value = payload.subarray(offset, Math.min(offset + 7, payload.length)) + offset += value.length + return { done: false, value } + }, + releaseLock() {}, + }), + }, { ifNoneMatch: '*', contentType: 'text/plain', metadata: { purpose: 'stream-test' } }) + await duringUpload + expect(requests).toEqual(['PUT', 'HEAD']) + expect(headers['content-length']).toBe(String(payload.length)) + expect(headers['transfer-encoding']).toBeUndefined() + expect(headers['if-none-match']).toBe('*') + expect(headers['content-type']).toBe('text/plain') + expect(headers['x-amz-meta-purpose']).toBe('stream-test') + expect(received).toEqual(payload) + expect(result.size).toBe(payload.length) + expect(await readdir(directory)).toEqual([]) + }) + + it('keeps streamed conditional failures atomic and removes the spool file', async () => { + const directory = await uploadTempDirectory() + let calls = 0 + let match: string | undefined + const { store } = await fixture((request, response) => { + calls++ + match = request.headers['if-match'] + request.resume() + request.on('end', () => xmlError(response, 412)) + }) + await expect(store.put('existing', bytesToObjectStream(new Uint8Array([1, 2])), { + ifMatchEtag: 'old-etag', + })).rejects.toMatchObject({ code: 'conflict' }) + expect(match).toBe('"old-etag"') + expect(calls).toBe(1) + expect(await readdir(directory)).toEqual([]) + }) + + it('uploads an empty stream with Content-Length zero', async () => { + const directory = await uploadTempDirectory() + let length: string | undefined + const { store } = await fixture((request, response) => { + if (request.method === 'PUT') { + length = request.headers['content-length'] + request.resume() + request.on('end', () => response.end()) + } else { + response.writeHead(200, { 'content-length': '0', 'etag': '"empty"' }) + response.end() + } + }) + expect((await store.put('empty', bytesToObjectStream(new Uint8Array()))).size).toBe(0) + expect(length).toBe('0') + expect(await readdir(directory)).toEqual([]) + }) + + it('does not send partial bytes or leak source errors when a stream fails', async () => { + const directory = await uploadTempDirectory() + let calls = 0 + const { store } = await fixture((_request, response) => { + calls++ + response.end() + }) + let reads = 0 + const cancel = vi.fn(async () => {}) + const releaseLock = vi.fn() + await expect(store.put('failed', { + getReader: () => ({ + read: async () => { + if (reads++) throw new Error('test-secret upstream-private-host') + return { done: false, value: new Uint8Array([1, 2, 3]) } + }, + cancel, + releaseLock, + }), + })).rejects.toMatchObject({ code: 'unavailable', message: 'S3 PUT failed' }) + expect(calls).toBe(0) + expect(cancel).toHaveBeenCalledOnce() + expect(releaseLock).toHaveBeenCalledOnce() + expect(await readdir(directory)).toEqual([]) + }) + + it('times out a pending source read, cancels it and removes temporary bytes', async () => { + const directory = await uploadTempDirectory() + let calls = 0 + const { store } = await fixture((_request, response) => { + calls++ + response.end() + }, { requestTimeoutMs: 100 }) + let reads = 0 + const cancel = vi.fn(async () => {}) + const releaseLock = vi.fn() + await expect(store.put('stalled', { + getReader: () => ({ + read: async () => reads++ + ? new Promise<{ done: boolean }>(() => {}) + : { done: false, value: new Uint8Array([1]) }, + cancel, + releaseLock, + }), + })).rejects.toMatchObject({ code: 'unavailable' }) + expect(calls).toBe(0) + expect(cancel).toHaveBeenCalledOnce() + expect(releaseLock).toHaveBeenCalledOnce() + expect(await readdir(directory)).toEqual([]) + }) + + it('uses the remaining upload deadline for HEAD and cleans up when it expires', async () => { + const directory = await uploadTempDirectory() + let headStarted = false + const { store } = await fixture((request, response) => { + if (request.method === 'PUT') { + request.resume() + request.on('end', () => response.end()) + } else { + headStarted = true + // This would finish within a fresh request timeout, but not within the + // time remaining after the slow source was consumed. + const timer = setTimeout(() => { + response.writeHead(200, { 'content-length': '1', 'etag': '"stored"' }) + response.end() + }, 200) + response.on('close', () => clearTimeout(timer)) + } + }, { requestTimeoutMs: 300 }) + let reads = 0 + await expect(store.put('head-timeout', { + getReader: () => ({ + read: async () => { + if (reads++) return { done: true } + await new Promise(resolve => setTimeout(resolve, 200)) + return { done: false, value: new Uint8Array([1]) } + }, + releaseLock() {}, + }), + })).rejects.toMatchObject({ code: 'unavailable' }) + expect(headStarted).toBe(true) + expect(await readdir(directory)).toEqual([]) }) it('bounds streamed downloads without buffering the whole response', async () => { diff --git a/packages/server/test/s3Probe.test.ts b/packages/server/test/s3Probe.test.ts new file mode 100644 index 00000000..af404d4c --- /dev/null +++ b/packages/server/test/s3Probe.test.ts @@ -0,0 +1,155 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { MemoryObjectStore, TBError } from '@tool-bridge/core' +import { createS3ObjectStore, type S3ObjectStore } from '../src/s3Objects' +import { probeS3ObjectStore } from '../src/s3Probe' + +vi.mock('../src/s3Objects', () => ({ createS3ObjectStore: vi.fn() })) + +const config = { + endpoint: 'https://objects.example.com', + bucket: 'probe-tests', + accessKeyId: 'test-access', + secretAccessKey: 'test-secret', +} + +beforeEach(() => vi.spyOn(console, 'warn').mockImplementation(() => {})) +afterEach(() => vi.restoreAllMocks()) + +function fixture() { + const memory = new MemoryObjectStore() + // Serialize the memory fixture's read-check-write sequence to model atomic S3 + // conditions; MemoryObjectStore itself awaits body conversion before writing. + let pending = Promise.resolve() + const store = { + head: vi.fn(memory.head.bind(memory)), + get: vi.fn(memory.get.bind(memory)), + list: vi.fn(memory.list.bind(memory)), + delete: vi.fn(memory.delete.bind(memory)), + close: vi.fn(), + put: vi.fn((...args) => { + const result = pending.then(() => memory.put(...args)) + pending = result.then(() => {}, () => {}) + return result + }), + } + vi.mocked(createS3ObjectStore).mockReturnValue(store) + return { store, memory } +} + +describe('S3 capability probe isolation and cleanup', () => { + it('checks all capabilities and cleans every attempted object', async () => { + const { store, memory } = fixture() + const result = await probeS3ObjectStore(config) + + expect(Object.values(result.checks).every(Boolean)).toBe(true) + expect(result.cleanupSucceeded).toBe(true) + expect((await memory.list('')).items).toEqual([]) + expect(store.list).toHaveBeenCalledTimes(3) + expect(store.close).toHaveBeenCalledOnce() + expect(console.warn).not.toHaveBeenCalled() + }) + + it.each(['before-write', 'after-write'] as const)( + 'a streaming error %s cannot cause a false pagination failure', + async (failurePoint) => { + const { store, memory } = fixture() + const put = store.put.getMockImplementation()! + store.put.mockImplementation(async (...args) => { + if (!args[0].endsWith('/stream')) return put(...args) + if (failurePoint === 'after-write') await put(...args) + throw new TBError('unavailable', 'S3 put failed') + }) + + const result = await probeS3ObjectStore(config) + + expect(result.checks.streaming).toBe(false) + expect(result.checks.pagination).toBe(true) + expect(result.cleanupSucceeded).toBe(true) + const attempted = new Set(store.put.mock.calls.map(([key]) => key)) + expect(new Set(store.delete.mock.calls.map(([key]) => key))).toEqual(attempted) + expect((await memory.list('')).items).toEqual([]) + expect(store.close).toHaveBeenCalledOnce() + }, + ) + + it('still rejects incomplete pagination', async () => { + const { store, memory } = fixture() + store.list.mockImplementation(async (prefix, options) => { + const result = await memory.list(prefix, options) + return { items: result.items } + }) + + const result = await probeS3ObjectStore(config) + + expect(result.checks.pagination).toBe(false) + expect(console.warn).toHaveBeenCalledExactlyOnceWith('S3 capability probe check failed', { + check: 'pagination', reason: 'unexpected_result', + }) + expect(result.cleanupSucceeded).toBe(true) + expect((await memory.list('')).items).toEqual([]) + }) + + it('reports cleanup failure while continuing cleanup of all attempted objects', async () => { + const { store, memory } = fixture() + store.delete.mockImplementation(async (key) => { + if (key.endsWith('/basic')) throw new TBError('permission_denied', 'S3 delete was denied') + await memory.delete(key) + }) + + const result = await probeS3ObjectStore(config) + + expect(Object.values(result.checks).every(Boolean)).toBe(true) + expect(result.cleanupSucceeded).toBe(false) + const attempted = new Set(store.put.mock.calls.map(([key]) => key)) + expect(new Set(store.delete.mock.calls.map(([key]) => key))).toEqual(attempted) + expect((await memory.list('')).items).toHaveLength(1) + expect(store.close).toHaveBeenCalledOnce() + }) + + it.each([ + 'S3 PUT failed (411)', + 'S3 HEAD failed (503)', + 'S3 GET was denied', + 'S3 LIST object not found', + 'S3 DELETE condition failed', + 'S3 GET deadline exceeded', + ])('retains the normalized operation and HTTP status: %s', async (message) => { + const { store } = fixture() + const put = store.put.getMockImplementation()! + store.put.mockImplementation(async (...args) => { + if (args[0].endsWith('/stream')) throw new TBError('unavailable', message) + return put(...args) + }) + + const result = await probeS3ObjectStore(config) + + expect(result.checks.streaming).toBe(false) + expect(console.warn).toHaveBeenCalledExactlyOnceWith('S3 capability probe check failed', { + check: 'streaming', reason: message, + }) + }) + + it.each([ + [new Error('test-secret upstream-private-host'), 'unexpected_error'], + [new TBError('unavailable', 'test-secret upstream-private-host'), 'operation_failed'], + [new TBError('unavailable', 'S3 PUT failed (411) test-secret'), 'operation_failed'], + [new TBError('unavailable', 'S3 PUT failed (411)\n'), 'operation_failed'], + [{ message: 'test-secret', endpoint: config.endpoint }, 'unexpected_error'], + ])('never logs arbitrary exceptions or forged TBError details: %#', async (error, reason) => { + const { store } = fixture() + const put = store.put.getMockImplementation()! + store.put.mockImplementation(async (...args) => { + if (args[0].endsWith('/stream')) throw error + return put(...args) + }) + + await probeS3ObjectStore(config) + + expect(console.warn).toHaveBeenCalledExactlyOnceWith('S3 capability probe check failed', { + check: 'streaming', reason, + }) + const logged = JSON.stringify(vi.mocked(console.warn).mock.calls) + for (const privateValue of [config.secretAccessKey, config.accessKeyId, config.endpoint, '__tool_bridge_internal__']) + expect(logged).not.toContain(privateValue) + }) +})