Add optional S3-based storage
This commit is contained in:
parent
6c812b44bb
commit
b93c840ee2
6 changed files with 135 additions and 46 deletions
21
src/db/db.ts
21
src/db/db.ts
|
|
@ -27,14 +27,23 @@ CREATE TABLE IF NOT EXISTS jobs (
|
||||||
num_files INTEGER DEFAULT 0,
|
num_files INTEGER DEFAULT 0,
|
||||||
FOREIGN KEY (user_id) REFERENCES users(id)
|
FOREIGN KEY (user_id) REFERENCES users(id)
|
||||||
);
|
);
|
||||||
PRAGMA user_version = 1;`);
|
CREATE TABLE IF NOT EXISTS storage_metadata (
|
||||||
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
|
user_id INTEGER NOT NULL,
|
||||||
|
job_id INTEGER NOT NULL,
|
||||||
|
file_name TEXT NOT NULL,
|
||||||
|
storage_key TEXT NOT NULL,
|
||||||
|
FOREIGN KEY (job_id) REFERENCES jobs(id),
|
||||||
|
FOREIGN KEY (user_id) REFERENCES users(id)
|
||||||
|
);
|
||||||
|
PRAGMA user_version = 2;`);
|
||||||
}
|
}
|
||||||
|
|
||||||
const dbVersion = (db.query("PRAGMA user_version").get() as { user_version?: number }).user_version;
|
const dbVersion = (db.query("PRAGMA user_version").get() as { user_version?: number }).user_version!;
|
||||||
if (dbVersion === 0) {
|
if (dbVersion < 2) {
|
||||||
db.exec("ALTER TABLE file_names ADD COLUMN status TEXT DEFAULT 'not started';");
|
db.exec("ALTER TABLE file_names ADD COLUMN storage_key TEXT;");
|
||||||
db.exec("PRAGMA user_version = 1;");
|
db.exec("PRAGMA user_version = 2;");
|
||||||
console.log("Updated database to version 1.");
|
console.log("Updated database to version 2.");
|
||||||
}
|
}
|
||||||
|
|
||||||
// enable WAL mode
|
// enable WAL mode
|
||||||
|
|
|
||||||
|
|
@ -1,18 +1,15 @@
|
||||||
import path from "node:path";
|
|
||||||
import { Elysia } from "elysia";
|
import { Elysia } from "elysia";
|
||||||
import sanitize from "sanitize-filename";
|
import sanitize from "sanitize-filename";
|
||||||
import * as tar from "tar";
|
|
||||||
import { outputDir } from "..";
|
|
||||||
import db from "../db/db";
|
import db from "../db/db";
|
||||||
import { WEBROOT } from "../helpers/env";
|
import { WEBROOT } from "../helpers/env";
|
||||||
import { userService } from "./user";
|
import { userService } from "./user";
|
||||||
|
import { getStorage } from "../storage/index";
|
||||||
|
|
||||||
export const download = new Elysia()
|
export const download = new Elysia()
|
||||||
.use(userService)
|
.use(userService)
|
||||||
.get(
|
.get(
|
||||||
"/download/:userId/:jobId/:fileName",
|
"/download/:userId/:jobId/:fileName",
|
||||||
async ({ params, redirect, user }) => {
|
async ({ params, redirect, user }) => {
|
||||||
const userId = user.id;
|
|
||||||
const job = await db
|
const job = await db
|
||||||
.query("SELECT * FROM jobs WHERE user_id = ? AND id = ?")
|
.query("SELECT * FROM jobs WHERE user_id = ? AND id = ?")
|
||||||
.get(user.id, params.jobId);
|
.get(user.id, params.jobId);
|
||||||
|
|
@ -20,44 +17,30 @@ export const download = new Elysia()
|
||||||
if (!job) {
|
if (!job) {
|
||||||
return redirect(`${WEBROOT}/results`, 302);
|
return redirect(`${WEBROOT}/results`, 302);
|
||||||
}
|
}
|
||||||
// parse from URL encoded string
|
|
||||||
const jobId = decodeURIComponent(params.jobId);
|
|
||||||
const fileName = sanitize(decodeURIComponent(params.fileName));
|
const fileName = sanitize(decodeURIComponent(params.fileName));
|
||||||
|
|
||||||
const filePath = `${outputDir}${userId}/${jobId}/${fileName}`;
|
const fileRow = db
|
||||||
return Bun.file(filePath);
|
.query(`
|
||||||
},
|
SELECT storage_key FROM file_names
|
||||||
{
|
WHERE job_id = ? AND file_name = ?
|
||||||
auth: true,
|
`,
|
||||||
},
|
|
||||||
)
|
)
|
||||||
.get(
|
.get(params.jobId, fileName) as { storage_key: string } | undefined;
|
||||||
"/archive/:jobId",
|
|
||||||
async ({ params, redirect, user }) => {
|
|
||||||
const userId = user.id;
|
|
||||||
const job = await db
|
|
||||||
.query("SELECT * FROM jobs WHERE user_id = ? AND id = ?")
|
|
||||||
.get(user.id, params.jobId);
|
|
||||||
|
|
||||||
if (!job) {
|
if (!fileRow) {
|
||||||
return redirect(`${WEBROOT}/results`, 302);
|
return redirect(`${WEBROOT}/results`, 302);
|
||||||
}
|
}
|
||||||
|
|
||||||
const jobId = decodeURIComponent(params.jobId);
|
const storage = getStorage();
|
||||||
const outputPath = `${outputDir}${userId}/${jobId}`;
|
const fileBuffer = await storage.get(fileRow.storage_key);
|
||||||
const outputTar = path.join(outputPath, `converted_files_${jobId}.tar`);
|
|
||||||
|
|
||||||
await tar.create(
|
return new Response(fileBuffer, {
|
||||||
{
|
headers: {
|
||||||
file: outputTar,
|
"Content-Type": "application/octet-stream",
|
||||||
cwd: outputPath,
|
"Content-Disposition": `attachment; filename="${fileName}"`,
|
||||||
filter: (path) => {
|
|
||||||
return !path.match(".*\\.tar");
|
|
||||||
},
|
},
|
||||||
},
|
});
|
||||||
["."],
|
|
||||||
);
|
|
||||||
return Bun.file(outputTar);
|
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
auth: true,
|
auth: true,
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,10 @@
|
||||||
import { Elysia, t } from "elysia";
|
import { Elysia, t } from "elysia";
|
||||||
import db from "../db/db";
|
import db from "../db/db";
|
||||||
import { WEBROOT } from "../helpers/env";
|
import { WEBROOT } from "../helpers/env";
|
||||||
import { uploadsDir } from "../index";
|
|
||||||
import { userService } from "./user";
|
import { userService } from "./user";
|
||||||
import sanitize from "sanitize-filename";
|
import sanitize from "sanitize-filename";
|
||||||
|
import { getStorage } from "../storage";
|
||||||
|
import crypto from "node:crypto";
|
||||||
|
|
||||||
export const upload = new Elysia().use(userService).post(
|
export const upload = new Elysia().use(userService).post(
|
||||||
"/upload",
|
"/upload",
|
||||||
|
|
@ -12,6 +13,8 @@ export const upload = new Elysia().use(userService).post(
|
||||||
return redirect(`${WEBROOT}/`, 302);
|
return redirect(`${WEBROOT}/`, 302);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const jobIdValue = jobId.value;
|
||||||
|
|
||||||
const existingJob = await db
|
const existingJob = await db
|
||||||
.query("SELECT * FROM jobs WHERE id = ? AND user_id = ?")
|
.query("SELECT * FROM jobs WHERE id = ? AND user_id = ?")
|
||||||
.get(jobId.value, user.id);
|
.get(jobId.value, user.id);
|
||||||
|
|
@ -20,17 +23,27 @@ export const upload = new Elysia().use(userService).post(
|
||||||
return redirect(`${WEBROOT}/`, 302);
|
return redirect(`${WEBROOT}/`, 302);
|
||||||
}
|
}
|
||||||
|
|
||||||
const userUploadsDir = `${uploadsDir}${user.id}/${jobId.value}/`;
|
const storage = getStorage();
|
||||||
|
|
||||||
|
const saveFile = async (file: File) => {
|
||||||
|
const sanitizedFileName = sanitize(file.name);
|
||||||
|
const storageKey = `${user.id}/${jobId.value}/${crypto.randomUUID()}`;
|
||||||
|
const buffer = Buffer.from(await file.arrayBuffer());
|
||||||
|
await storage.save(storageKey, buffer);
|
||||||
|
|
||||||
|
db.query(`
|
||||||
|
INSERT INTO file_names (job_id, file_name, storage_key)
|
||||||
|
VALUES (?, ?, ?)
|
||||||
|
`).run(jobIdValue, sanitizedFileName, storageKey);
|
||||||
|
};
|
||||||
|
|
||||||
if (body?.file) {
|
if (body?.file) {
|
||||||
if (Array.isArray(body.file)) {
|
if (Array.isArray(body.file)) {
|
||||||
for (const file of body.file) {
|
for (const file of body.file) {
|
||||||
const santizedFileName = sanitize(file.name);
|
await saveFile(file);
|
||||||
await Bun.write(`${userUploadsDir}${santizedFileName}`, file);
|
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
const santizedFileName = sanitize(body.file["name"]);
|
await saveFile(body.file);;
|
||||||
await Bun.write(`${userUploadsDir}${santizedFileName}`, body.file);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
28
src/storage/LocalStorageAdapter.ts
Normal file
28
src/storage/LocalStorageAdapter.ts
Normal file
|
|
@ -0,0 +1,28 @@
|
||||||
|
import { IStorageAdapter } from "./index";
|
||||||
|
import { promises as fs } from "fs";
|
||||||
|
import path from "path";
|
||||||
|
|
||||||
|
export class LocalStorageAdapter implements IStorageAdapter {
|
||||||
|
baseDir: string;
|
||||||
|
|
||||||
|
constructor(baseDir: string) {
|
||||||
|
this.baseDir = baseDir;
|
||||||
|
}
|
||||||
|
|
||||||
|
async save(key: string, data: Buffer): Promise<string> {
|
||||||
|
const fullPath = path.join(this.baseDir, key);
|
||||||
|
await fs.mkdir(path.dirname(fullPath), { recursive: true });
|
||||||
|
await fs.writeFile(fullPath, data);
|
||||||
|
return key;
|
||||||
|
}
|
||||||
|
|
||||||
|
async get(key: string): Promise<Buffer> {
|
||||||
|
const fullPath = path.join(this.baseDir, key);
|
||||||
|
return fs.readFile(fullPath);
|
||||||
|
}
|
||||||
|
|
||||||
|
async delete(key: string): Promise<void> {
|
||||||
|
const fullPath = path.join(this.baseDir, key);
|
||||||
|
await fs.unlink(fullPath);
|
||||||
|
}
|
||||||
|
}
|
||||||
36
src/storage/S3StorageAdapter.ts
Normal file
36
src/storage/S3StorageAdapter.ts
Normal file
|
|
@ -0,0 +1,36 @@
|
||||||
|
import { s3, S3File } from "bun";
|
||||||
|
import { IStorageAdapter } from "./index";
|
||||||
|
|
||||||
|
export class S3StorageAdapter implements IStorageAdapter {
|
||||||
|
private bucket: string;
|
||||||
|
|
||||||
|
constructor(bucket: string) {
|
||||||
|
this.bucket = bucket;
|
||||||
|
}
|
||||||
|
|
||||||
|
async save(key: string, data: Buffer): Promise<string> {
|
||||||
|
const file: S3File = s3.file(key, {
|
||||||
|
bucket: this.bucket,
|
||||||
|
acl: "private",
|
||||||
|
});
|
||||||
|
|
||||||
|
await file.write(data);
|
||||||
|
return key;
|
||||||
|
}
|
||||||
|
|
||||||
|
async get(key: string): Promise<Buffer> {
|
||||||
|
const file: S3File = s3.file(key, {
|
||||||
|
bucket: this.bucket,
|
||||||
|
});
|
||||||
|
|
||||||
|
return Buffer.from(await file.bytes());
|
||||||
|
}
|
||||||
|
|
||||||
|
async delete(key: string): Promise<void> {
|
||||||
|
const file: S3File = s3.file(key, {
|
||||||
|
bucket: this.bucket,
|
||||||
|
})
|
||||||
|
|
||||||
|
await file.delete();
|
||||||
|
}
|
||||||
|
}
|
||||||
20
src/storage/index.ts
Normal file
20
src/storage/index.ts
Normal file
|
|
@ -0,0 +1,20 @@
|
||||||
|
import { LocalStorageAdapter } from "./LocalStorageAdapter";
|
||||||
|
import { S3StorageAdapter } from "./S3StorageAdapter";
|
||||||
|
|
||||||
|
export interface IStorageAdapter {
|
||||||
|
save(key: string, data: Buffer): Promise<string>;
|
||||||
|
get(key: string): Promise<Buffer>;
|
||||||
|
delete(key: string): Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function getStorage(): IStorageAdapter {
|
||||||
|
if (process.env.STORAGE_BACKEND === "s3") {
|
||||||
|
if (!process.env.S3_BUCKET) {
|
||||||
|
throw new Error("S3_BUCKET must be set when STORAGE_BACKEND=s3");
|
||||||
|
}
|
||||||
|
|
||||||
|
return new S3StorageAdapter(process.env.S3_BUCKET);
|
||||||
|
}
|
||||||
|
|
||||||
|
return new LocalStorageAdapter("./data");
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue