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.
This commit is contained in:
@@ -115,6 +115,9 @@ const startServer = async () => {
|
|||||||
await seedDatabase({ disconnectOnComplete: false });
|
await seedDatabase({ disconnectOnComplete: false });
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const { ensureAppBuckets } = require('./utils/s3Client');
|
||||||
|
await ensureAppBuckets();
|
||||||
|
|
||||||
const registerEventListeners = require('./events/eventListeners');
|
const registerEventListeners = require('./events/eventListeners');
|
||||||
const { startPaymentReminderJob } = require('./jobs/paymentReminderJob');
|
const { startPaymentReminderJob } = require('./jobs/paymentReminderJob');
|
||||||
const { startNotificationRetryJob } = require('./jobs/notificationRetryJob');
|
const { startNotificationRetryJob } = require('./jobs/notificationRetryJob');
|
||||||
|
|||||||
@@ -144,6 +144,11 @@
|
|||||||
"en": "File is required.",
|
"en": "File is required.",
|
||||||
"fa": "ارسال فایل الزامی است."
|
"fa": "ارسال فایل الزامی است."
|
||||||
},
|
},
|
||||||
|
"FILE_COMMIT_FAILED": {
|
||||||
|
"statusCode": 502,
|
||||||
|
"en": "Failed to store the uploaded file. Please try again.",
|
||||||
|
"fa": "ذخیره فایل آپلودشده ناموفق بود. لطفا دوباره تلاش کنید."
|
||||||
|
},
|
||||||
"SEED_LOCKED": {
|
"SEED_LOCKED": {
|
||||||
"statusCode": 409,
|
"statusCode": 409,
|
||||||
"en": "Database seeding has already been completed and is locked.",
|
"en": "Database seeding has already been completed and is locked.",
|
||||||
|
|||||||
+107
-20
@@ -4,13 +4,15 @@ const {
|
|||||||
S3Client,
|
S3Client,
|
||||||
PutObjectCommand,
|
PutObjectCommand,
|
||||||
GetObjectCommand,
|
GetObjectCommand,
|
||||||
CopyObjectCommand,
|
|
||||||
DeleteObjectCommand,
|
DeleteObjectCommand,
|
||||||
ListObjectsV2Command
|
ListObjectsV2Command,
|
||||||
|
HeadBucketCommand,
|
||||||
|
CreateBucketCommand
|
||||||
} = require('@aws-sdk/client-s3');
|
} = require('@aws-sdk/client-s3');
|
||||||
const { getSignedUrl } = require('@aws-sdk/s3-request-presigner');
|
const { getSignedUrl } = require('@aws-sdk/s3-request-presigner');
|
||||||
const config = require('../config/config');
|
const config = require('../config/config');
|
||||||
const logger = require('./logger');
|
const logger = require('./logger');
|
||||||
|
const AppError = require('./AppError');
|
||||||
|
|
||||||
let s3ClientInstance = null;
|
let s3ClientInstance = null;
|
||||||
|
|
||||||
@@ -64,20 +66,31 @@ const getS3Client = () => {
|
|||||||
return s3ClientInstance;
|
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 logS3Error = (label, error) => {
|
||||||
const status = error?.$metadata?.httpStatusCode;
|
const status = error?.$metadata?.httpStatusCode;
|
||||||
const raw = error?.$response;
|
const raw = error?.$response;
|
||||||
let bodyPreview = '';
|
let bodyPreview = '';
|
||||||
try {
|
try {
|
||||||
const body = raw?.body || raw?.reason || '';
|
bodyPreview = previewErrorBody(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);
|
|
||||||
} catch {
|
} catch {
|
||||||
bodyPreview = '';
|
bodyPreview = '';
|
||||||
}
|
}
|
||||||
logger.error(
|
logger.error(
|
||||||
`[S3 Storage ERROR] ${label}: ${error.message}`
|
`[S3 Storage ERROR] ${label}: ${error.message}`
|
||||||
+ (status ? ` (HTTP ${status})` : '')
|
+ (status ? ` (HTTP ${status})` : '')
|
||||||
|
+ (error?.name ? ` name=${error.name}` : '')
|
||||||
|
+ (error?.Code ? ` code=${error.Code}` : '')
|
||||||
+ (bodyPreview ? ` body=${bodyPreview.replace(/\s+/g, ' ')}` : '')
|
+ (bodyPreview ? ` body=${bodyPreview.replace(/\s+/g, ' ')}` : '')
|
||||||
);
|
);
|
||||||
if (status === 307 || /Temporary Redirect/i.test(String(error.message) + bodyPreview)) {
|
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') => {
|
const uploadToTempBucket = async (fileBuffer, filename, contentType = 'application/octet-stream') => {
|
||||||
try {
|
try {
|
||||||
|
await ensureBucket(config.S3_TEMP_BUCKET);
|
||||||
const client = getS3Client();
|
const client = getS3Client();
|
||||||
const command = new PutObjectCommand({
|
const command = new PutObjectCommand({
|
||||||
Bucket: config.S3_TEMP_BUCKET,
|
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.
|
* Move a temp object into a target bucket, then remove the temp object.
|
||||||
* @param {string} tempFilename
|
* SeaweedFS CopyObject returns InternalError 500 across buckets, so this
|
||||||
* @param {string} destinationKey
|
* rewrites via GetObject + PutObject (PutObject already works for temp uploads).
|
||||||
* @param {string} [targetBucketKind='certificates'] - 'certificates' | 'documents' | bucket name
|
|
||||||
*/
|
*/
|
||||||
const commitTempFile = async (tempFilename, destinationKey = null, targetBucketKind = 'certificates') => {
|
const commitTempFile = async (tempFilename, destinationKey = null, targetBucketKind = 'certificates') => {
|
||||||
const targetKey = destinationKey || tempFilename;
|
const targetKey = destinationKey || tempFilename;
|
||||||
const targetBucket = resolveBucket(targetBucketKind);
|
const targetBucket = resolveBucket(targetBucketKind);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
await ensureBucket(targetBucket);
|
||||||
const client = getS3Client();
|
const client = getS3Client();
|
||||||
|
await rewriteTempObject(client, tempFilename, targetBucket, targetKey);
|
||||||
|
|
||||||
const copyCommand = new CopyObjectCommand({
|
try {
|
||||||
CopySource: `${config.S3_TEMP_BUCKET}/${tempFilename}`,
|
await client.send(new DeleteObjectCommand({
|
||||||
Bucket: targetBucket,
|
|
||||||
Key: targetKey
|
|
||||||
});
|
|
||||||
await client.send(copyCommand);
|
|
||||||
|
|
||||||
const deleteCommand = new DeleteObjectCommand({
|
|
||||||
Bucket: config.S3_TEMP_BUCKET,
|
Bucket: config.S3_TEMP_BUCKET,
|
||||||
Key: tempFilename
|
Key: tempFilename
|
||||||
});
|
}));
|
||||||
await client.send(deleteCommand);
|
} catch (deleteError) {
|
||||||
|
logger.warn(`[S3 Storage] Committed ${targetKey} but failed to delete temp ${tempFilename}: ${deleteError.message}`);
|
||||||
|
}
|
||||||
|
|
||||||
const fileUrl = buildPublicUrl(targetBucket, targetKey);
|
const fileUrl = buildPublicUrl(targetBucket, targetKey);
|
||||||
logger.info(`[S3 Storage] Committed ${tempFilename} → ${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)
|
fileUrl: buildPublicUrl(targetBucket, targetKey)
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
throw error;
|
throw new AppError('FILE_COMMIT_FAILED');
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -234,5 +319,7 @@ module.exports = {
|
|||||||
resolveBucket,
|
resolveBucket,
|
||||||
isBucketPublic,
|
isBucketPublic,
|
||||||
buildPublicUrl,
|
buildPublicUrl,
|
||||||
|
ensureBucket,
|
||||||
|
ensureAppBuckets,
|
||||||
BUCKETS
|
BUCKETS
|
||||||
};
|
};
|
||||||
|
|||||||
Reference in New Issue
Block a user