From 23edffba6afe1e4417db751c88b6c53161d7837b Mon Sep 17 00:00:00 2001 From: Kavehhn174 Date: Sat, 15 Aug 2026 03:33:55 +0330 Subject: [PATCH] fix: rewrite temp files with PutObject instead of SeaweedFS CopyObject CopyObject returns InternalError 500 across buckets, so commits now Get+Put and ensure the documents bucket exists on startup. --- app.js | 3 ++ utils/errors.json | 5 ++ utils/s3Client.js | 131 ++++++++++++++++++++++++++++++++++++++-------- 3 files changed, 117 insertions(+), 22 deletions(-) diff --git a/app.js b/app.js index 7763581..b6a98e2 100644 --- a/app.js +++ b/app.js @@ -115,6 +115,9 @@ const startServer = async () => { await seedDatabase({ disconnectOnComplete: false }); } + const { ensureAppBuckets } = require('./utils/s3Client'); + await ensureAppBuckets(); + const registerEventListeners = require('./events/eventListeners'); const { startPaymentReminderJob } = require('./jobs/paymentReminderJob'); const { startNotificationRetryJob } = require('./jobs/notificationRetryJob'); diff --git a/utils/errors.json b/utils/errors.json index f191fab..8de852b 100644 --- a/utils/errors.json +++ b/utils/errors.json @@ -144,6 +144,11 @@ "en": "File is required.", "fa": "ارسال فایل الزامی است." }, + "FILE_COMMIT_FAILED": { + "statusCode": 502, + "en": "Failed to store the uploaded file. Please try again.", + "fa": "ذخیره فایل آپلودشده ناموفق بود. لطفا دوباره تلاش کنید." + }, "SEED_LOCKED": { "statusCode": 409, "en": "Database seeding has already been completed and is locked.", diff --git a/utils/s3Client.js b/utils/s3Client.js index 1bae207..0c03324 100644 --- a/utils/s3Client.js +++ b/utils/s3Client.js @@ -4,13 +4,15 @@ const { S3Client, PutObjectCommand, GetObjectCommand, - CopyObjectCommand, DeleteObjectCommand, - ListObjectsV2Command + ListObjectsV2Command, + HeadBucketCommand, + CreateBucketCommand } = require('@aws-sdk/client-s3'); const { getSignedUrl } = require('@aws-sdk/s3-request-presigner'); const config = require('../config/config'); const logger = require('./logger'); +const AppError = require('./AppError'); let s3ClientInstance = null; @@ -64,20 +66,31 @@ const getS3Client = () => { return s3ClientInstance; }; +const previewErrorBody = (body) => { + if (!body) return ''; + if (typeof body === 'string') return body.slice(0, 300); + if (Buffer.isBuffer(body)) return body.toString('utf8').slice(0, 300); + try { + return JSON.stringify(body).slice(0, 300); + } catch { + return Object.prototype.toString.call(body); + } +}; + const logS3Error = (label, error) => { const status = error?.$metadata?.httpStatusCode; const raw = error?.$response; let bodyPreview = ''; try { - const body = raw?.body || raw?.reason || ''; - if (typeof body === 'string') bodyPreview = body.slice(0, 300); - else if (body && typeof body.toString === 'function') bodyPreview = String(body).slice(0, 300); + bodyPreview = previewErrorBody(raw?.body || raw?.reason || ''); } catch { bodyPreview = ''; } logger.error( `[S3 Storage ERROR] ${label}: ${error.message}` + (status ? ` (HTTP ${status})` : '') + + (error?.name ? ` name=${error.name}` : '') + + (error?.Code ? ` code=${error.Code}` : '') + (bodyPreview ? ` body=${bodyPreview.replace(/\s+/g, ' ')}` : '') ); if (status === 307 || /Temporary Redirect/i.test(String(error.message) + bodyPreview)) { @@ -88,8 +101,69 @@ const logS3Error = (label, error) => { } }; +const streamToBuffer = async (body) => { + if (!body) return Buffer.alloc(0); + if (Buffer.isBuffer(body)) return body; + if (typeof body.transformToByteArray === 'function') { + return Buffer.from(await body.transformToByteArray()); + } + const chunks = []; + for await (const chunk of body) { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + } + return Buffer.concat(chunks); +}; + +const ensuredBuckets = new Set(); + +const ensureBucket = async (bucketNameOrKind) => { + const bucket = resolveBucket(bucketNameOrKind); + if (ensuredBuckets.has(bucket)) return bucket; + + const client = getS3Client(); + try { + await client.send(new HeadBucketCommand({ Bucket: bucket })); + ensuredBuckets.add(bucket); + return bucket; + } catch (headErr) { + logS3Error(`HeadBucket ${bucket} (will try create)`, headErr); + } + + try { + await client.send(new CreateBucketCommand({ Bucket: bucket })); + logger.info(`[S3 Storage] Created bucket ${bucket}`); + } catch (createErr) { + const code = createErr?.name || createErr?.Code || ''; + if (code !== 'BucketAlreadyOwnedByYou' && code !== 'BucketAlreadyExists') { + logS3Error(`Failed to create bucket ${bucket}`, createErr); + throw createErr; + } + } + + ensuredBuckets.add(bucket); + return bucket; +}; + +const ensureAppBuckets = async () => { + const names = [ + config.S3_TEMP_BUCKET, + config.S3_CERTIFICATES_BUCKET, + config.S3_DOCUMENTS_BUCKET + ].filter(Boolean); + + for (const name of [...new Set(names)]) { + try { + await ensureBucket(name); + logger.info(`[S3 Storage] Bucket ready: ${name}`); + } catch (err) { + logger.error(`[S3 Storage] Could not ensure bucket ${name}: ${err.message}`); + } + } +}; + const uploadToTempBucket = async (fileBuffer, filename, contentType = 'application/octet-stream') => { try { + await ensureBucket(config.S3_TEMP_BUCKET); const client = getS3Client(); const command = new PutObjectCommand({ Bucket: config.S3_TEMP_BUCKET, @@ -112,31 +186,42 @@ const uploadToTempBucket = async (fileBuffer, filename, contentType = 'applicati } }; +const rewriteTempObject = async (client, tempFilename, targetBucket, targetKey) => { + const source = await client.send(new GetObjectCommand({ + Bucket: config.S3_TEMP_BUCKET, + Key: tempFilename + })); + const body = await streamToBuffer(source.Body); + await client.send(new PutObjectCommand({ + Bucket: targetBucket, + Key: targetKey, + Body: body, + ContentType: source.ContentType || 'application/octet-stream' + })); +}; + /** - * Copy a temp object into a target private/public bucket, then remove the temp object. - * @param {string} tempFilename - * @param {string} destinationKey - * @param {string} [targetBucketKind='certificates'] - 'certificates' | 'documents' | bucket name + * Move a temp object into a target bucket, then remove the temp object. + * SeaweedFS CopyObject returns InternalError 500 across buckets, so this + * rewrites via GetObject + PutObject (PutObject already works for temp uploads). */ const commitTempFile = async (tempFilename, destinationKey = null, targetBucketKind = 'certificates') => { const targetKey = destinationKey || tempFilename; const targetBucket = resolveBucket(targetBucketKind); try { + await ensureBucket(targetBucket); const client = getS3Client(); + await rewriteTempObject(client, tempFilename, targetBucket, targetKey); - const copyCommand = new CopyObjectCommand({ - CopySource: `${config.S3_TEMP_BUCKET}/${tempFilename}`, - Bucket: targetBucket, - Key: targetKey - }); - await client.send(copyCommand); - - const deleteCommand = new DeleteObjectCommand({ - Bucket: config.S3_TEMP_BUCKET, - Key: tempFilename - }); - await client.send(deleteCommand); + try { + await client.send(new DeleteObjectCommand({ + Bucket: config.S3_TEMP_BUCKET, + Key: tempFilename + })); + } catch (deleteError) { + logger.warn(`[S3 Storage] Committed ${targetKey} but failed to delete temp ${tempFilename}: ${deleteError.message}`); + } const fileUrl = buildPublicUrl(targetBucket, targetKey); logger.info(`[S3 Storage] Committed ${tempFilename} → ${targetBucket}/${targetKey}`); @@ -151,7 +236,7 @@ const commitTempFile = async (tempFilename, destinationKey = null, targetBucketK fileUrl: buildPublicUrl(targetBucket, targetKey) }; } - throw error; + throw new AppError('FILE_COMMIT_FAILED'); } }; @@ -234,5 +319,7 @@ module.exports = { resolveBucket, isBucketPublic, buildPublicUrl, + ensureBucket, + ensureAppBuckets, BUCKETS };