diff --git a/package.json b/package.json index a2285ec..357634e 100644 --- a/package.json +++ b/package.json @@ -9,12 +9,14 @@ "mcp": "bun run src/mcp.ts" }, "devDependencies": { - "@types/bun": "^1.3.1" + "@types/bun": "^1.3.1", + "@types/nodemailer": "^6.4.17" }, "peerDependencies": { "typescript": "^5.0.0" }, "dependencies": { + "@aws-sdk/client-s3": "^3.700.0", "@elysiajs/static": "^1.4.0", "@modelcontextprotocol/sdk": "^1.30.0", "@types/crypto-js": "^4.2.2", @@ -31,6 +33,7 @@ "isomorphic-dompurify": "^2.18.0", "js-yaml": "^4.1.1", "kysely": "^0.27.4", + "nodemailer": "^6.9.16", "pg": "^8.13.1", "rolling-rate-limiter": "^0.4.2" } diff --git a/src/index.ts b/src/index.ts index 1075c8a..851ed0e 100644 --- a/src/index.ts +++ b/src/index.ts @@ -8,6 +8,8 @@ import video from '@/router/video' import websocket from '@/router/websocket' import html from '@/router/html' import mcp from '@/router/mcp' +import admin from '@/router/admin' +import { startRestorePoller } from '@/utils/glacierPoller' const app = new Elysia() app.use(latest) @@ -17,6 +19,11 @@ app.use(video) app.use(websocket) app.use(html) app.use(mcp) +app.use(admin) + +if (process.env.IS_PRIMARY === 'true') { + startRestorePoller() +} app.onRequest(({ set, request }: { set: { headers: Record }, request: Request }) => { set.headers['Onion-Location'] = 'http://tubey5btlzxkcjpxpj2c7irrbhvgu3noouobndafuhbw4i5ndvn4v7qd.onion/' + request.url.split('/').at(-1) }) diff --git a/src/router/admin.ts b/src/router/admin.ts new file mode 100644 index 0000000..b50099d --- /dev/null +++ b/src/router/admin.ts @@ -0,0 +1,112 @@ +import { Elysia, t } from 'elysia' +import { db } from '@/utils/database' +import { m, eta, error as errorPage } from '@/utils/html' +import { + getSessionToken, + createAdminSession, + isValidAdminSession, + destroyAdminSession, + buildSessionCookie, + clearSessionCookie +} from '@/utils/adminAuth' +import { initiateRestore } from '@/utils/glacier' + +const app = new Elysia({ prefix: '/admin' }) + +app.onBeforeHandle(async ({ path, request, set, headers, redirect }) => { + if (path === '/admin/login') return + + const token = getSessionToken(headers.cookie) + if (await isValidAdminSession(token)) return + + if (request.method === 'GET') return redirect('/admin/login') + + set.status = 401 + return { success: false, message: 'Unauthorized' } +}) + +app.get('/login', async ({ set }) => { + set.headers['Content-Type'] = 'text/html; charset=utf-8' + return await m(eta.render('./admin/login', { + title: 'Admin Login | PreserveTube' + })) +}) + +app.post('/login', async ({ body, set, redirect }) => { + if (body.password !== process.env.ADMIN_SECRET) { + set.headers['Content-Type'] = 'text/html; charset=utf-8' + set.status = 401 + return await m(eta.render('./admin/login', { + title: 'Admin Login | PreserveTube', + loginError: 'Incorrect password.' + })) + } + + const token = await createAdminSession() + set.headers['Set-Cookie'] = buildSessionCookie(token) + return redirect('/admin') +}, { + body: t.Object({ + password: t.String() + }) +}) + +app.post('/logout', async ({ headers, redirect }) => { + const token = getSessionToken(headers.cookie) + if (token) await destroyAdminSession(token) + return new Response(null, { + status: 302, + headers: { 'Set-Cookie': clearSessionCookie(), Location: '/admin/login' } + }) +}) + +app.get('/', async ({ set }) => { + const requests = await db.selectFrom('restore_requests') + .selectAll() + .orderBy('created_at', 'desc') + .execute() + + set.headers['Content-Type'] = 'text/html; charset=utf-8' + return await m(eta.render('./admin/dashboard', { + title: 'Admin | PreserveTube', + requests + })) +}) + +app.post('/restore', async ({ body, redirect, error }) => { + const { videoId, requesterEmail } = body + + const video = await db.selectFrom('videos') + .select(['id', 'deletion_stage']) + .where('id', '=', videoId) + .executeTakeFirst() + + if (!video) return error(404, 'No archived video found with that ID.') + if (video.deletion_stage !== 'cold_storage') return error(400, 'That video is not currently in cold storage.') + + const inserted = await db.insertInto('restore_requests') + .values({ + videoId, + requester_email: requesterEmail, + status: 'requested' + }) + .returning('uuid') + .executeTakeFirstOrThrow() + + await initiateRestore(videoId) + + await db.updateTable('restore_requests') + .set({ status: 'restoring', aws_restore_requested_at: new Date(), updated_at: new Date() }) + .where('uuid', '=', inserted.uuid) + .execute() + + return redirect('/admin') +}, { + body: t.Object({ + videoId: t.String(), + requesterEmail: t.String() + }) +}) + +app.onError(errorPage) +export default app diff --git a/src/templates/admin/dashboard.eta b/src/templates/admin/dashboard.eta new file mode 100644 index 0000000..a7f9d3c --- /dev/null +++ b/src/templates/admin/dashboard.eta @@ -0,0 +1,130 @@ +<% layout('../layout') %> + +
+
+

Admin

+
+ +
+
+ +

Glacier Restore

+

Restore a video from Glacier cold storage back to hot storage. AWS retrieval takes up to 48 hours; once it's done, the video is automatically re-uploaded and the requester is emailed.

+ +
+ + + +
+ + + + + + + + + + + + + <% it.requests.forEach(function(r){ %> + + + + + + + + <% }) %> + <% if (it.requests.length === 0) { %> + + <% } %> + +
VideoRequesterStatusRequestedError
<%= r.videoId %><%= r.requester_email %><%= r.status %><%= r.created_at ? new Date(r.created_at).toISOString() : '' %><%= r.error_message || '' %>
No restore requests yet.
+
+ + diff --git a/src/templates/admin/login.eta b/src/templates/admin/login.eta new file mode 100644 index 0000000..5048ecc --- /dev/null +++ b/src/templates/admin/login.eta @@ -0,0 +1,53 @@ +<% layout('../layout') %> + +
+

Admin Login

+ <% if (it.loginError) { %> +

<%= it.loginError %>

+ <% } %> +
+ + +
+
+ + diff --git a/src/types.ts b/src/types.ts index 1ec629c..6f83625 100644 --- a/src/types.ts +++ b/src/types.ts @@ -9,6 +9,7 @@ export interface Database { videos: VideosTable reports: ReportsTable files: FilesTable + restore_requests: RestoreRequestsTable } export interface VideosTable { @@ -64,3 +65,19 @@ export interface FilesTable { export type File = Selectable export type NewFile = Insertable + +export interface RestoreRequestsTable { + uuid: Generated + videoId: string + requester_email: string + status: 'requested' | 'restoring' | 'restored' | 'reuploading' | 'reuploaded' | 'failed' + aws_restore_requested_at: Date | null + aws_restore_expiry: Date | null + error_message: string | null + created_at: Generated + updated_at: Generated +} + +export type RestoreRequest = Selectable +export type NewRestoreRequest = Insertable +export type UpdateRestoreRequest = Updateable diff --git a/src/utils/adminAuth.ts b/src/utils/adminAuth.ts new file mode 100644 index 0000000..241a78d --- /dev/null +++ b/src/utils/adminAuth.ts @@ -0,0 +1,52 @@ +import redis from '@/utils/redis' + +const SESSION_COOKIE = 'pt_admin_session' +const SESSION_TTL_SECONDS = 60 * 60 * 12 + +const parseCookieHeader = (cookieHeader?: string): Record => { + if (!cookieHeader) return {} + + return cookieHeader.split(';').reduce>((acc, part) => { + const [rawKey, ...rawValue] = part.trim().split('=') + if (!rawKey || rawValue.length === 0) return acc + + acc[rawKey] = decodeURIComponent(rawValue.join('=')) + return acc + }, {}) +} + +function getSessionToken(cookieHeader?: string): string | undefined { + return parseCookieHeader(cookieHeader)[SESSION_COOKIE] +} + +async function createAdminSession(): Promise { + const token = crypto.randomUUID() + await redis.set(`admin:session:${token}`, '1', 'EX', SESSION_TTL_SECONDS) + return token +} + +async function isValidAdminSession(token: string | undefined): Promise { + if (!token) return false + return (await redis.get(`admin:session:${token}`)) === '1' +} + +async function destroyAdminSession(token: string): Promise { + await redis.del(`admin:session:${token}`) +} + +function buildSessionCookie(token: string): string { + return `${SESSION_COOKIE}=${encodeURIComponent(token)}; Max-Age=${SESSION_TTL_SECONDS}; Path=/admin; HttpOnly; SameSite=Strict; Secure` +} + +function clearSessionCookie(): string { + return `${SESSION_COOKIE}=; Max-Age=0; Path=/admin; HttpOnly; SameSite=Strict; Secure` +} + +export { + getSessionToken, + createAdminSession, + isValidAdminSession, + destroyAdminSession, + buildSessionCookie, + clearSessionCookie +} diff --git a/src/utils/glacier.ts b/src/utils/glacier.ts new file mode 100644 index 0000000..38ee54c --- /dev/null +++ b/src/utils/glacier.ts @@ -0,0 +1,60 @@ +import * as fs from 'node:fs' +import { Readable } from 'node:stream' +import { pipeline } from 'node:stream/promises' +import { S3Client, RestoreObjectCommand, HeadObjectCommand, GetObjectCommand } from '@aws-sdk/client-s3' + +const s3 = new S3Client({ region: process.env.AWS_REGION }) + +function objectKey(videoId: string) { + return `${videoId}.mp4` +} + +async function initiateRestore(videoId: string, days = 5) { + try { + await s3.send(new RestoreObjectCommand({ + Bucket: process.env.GLACIER_BUCKET, + Key: objectKey(videoId), + RestoreRequest: { + Days: days, + GlacierJobParameters: { Tier: 'Standard' } + } + })) + } catch (error: unknown) { + const err = error as { name?: string } + if (err.name === 'RestoreAlreadyInProgress') return + throw error + } +} + +async function checkRestoreStatus(videoId: string): Promise<{ ongoing: boolean, expiry: Date | null }> { + const head = await s3.send(new HeadObjectCommand({ + Bucket: process.env.GLACIER_BUCKET, + Key: objectKey(videoId) + })) + + const restoreHeader = head.Restore + if (!restoreHeader) return { ongoing: true, expiry: null } + + const ongoing = /ongoing-request="true"/.test(restoreHeader) + const expiryMatch = restoreHeader.match(/expiry-date="([^"]+)"/) + const expiry = expiryMatch ? new Date(expiryMatch[1]) : null + + return { ongoing, expiry } +} + +async function downloadRestoredObjectToFile(videoId: string, destPath: string) { + const object = await s3.send(new GetObjectCommand({ + Bucket: process.env.GLACIER_BUCKET, + Key: objectKey(videoId) + })) + + if (!object.Body) throw new Error(`No body returned for restored object ${objectKey(videoId)}`) + + const nodeStream = object.Body instanceof Readable + ? object.Body + : Readable.fromWeb(object.Body as any) + + await pipeline(nodeStream, fs.createWriteStream(destPath)) +} + +export { initiateRestore, checkRestoreStatus, downloadRestoredObjectToFile } diff --git a/src/utils/glacierPoller.ts b/src/utils/glacierPoller.ts new file mode 100644 index 0000000..dc48b5a --- /dev/null +++ b/src/utils/glacierPoller.ts @@ -0,0 +1,74 @@ +import * as fs from 'node:fs' +import { db } from '@/utils/database' +import redis from '@/utils/redis' +import { uploadVideo } from '@/utils/upload' +import { checkRestoreStatus, downloadRestoredObjectToFile } from '@/utils/glacier' +import { sendRestoreCompleteEmail } from '@/utils/mail' + +const POLL_INTERVAL_MS = 10 * 60000 + +async function processRestoringRow(row: { uuid: string, videoId: string, requester_email: string }) { + const { ongoing, expiry } = await checkRestoreStatus(row.videoId) + if (ongoing) return + + await db.updateTable('restore_requests') + .set({ status: 'reuploading', aws_restore_expiry: expiry, updated_at: new Date() }) + .where('uuid', '=', row.uuid) + .execute() + + const filePath = `./videos/${row.videoId}.mp4` + try { + await downloadRestoredObjectToFile(row.videoId, filePath) + const videoUrl = await uploadVideo(filePath) + + await db.updateTable('videos') + .set({ deletion_stage: null, source: videoUrl }) + .where('id', '=', row.videoId) + .execute() + + await redis.del(`watch:${row.videoId}:html`) + await redis.del('deletion:html') + + await db.updateTable('restore_requests') + .set({ status: 'reuploaded', updated_at: new Date() }) + .where('uuid', '=', row.uuid) + .execute() + + const video = await db.selectFrom('videos') + .select(['title']) + .where('id', '=', row.videoId) + .executeTakeFirst() + + await sendRestoreCompleteEmail(row.requester_email, row.videoId, video?.title) + } finally { + if (fs.existsSync(filePath)) fs.unlinkSync(filePath) + } +} + +async function pollRestores() { + const restoringRows = await db.selectFrom('restore_requests') + .select(['uuid', 'videoId', 'requester_email']) + .where('status', '=', 'restoring') + .execute() + + for (const row of restoringRows) { + try { + await processRestoringRow(row) + } catch (error: unknown) { + const err = error as Error + console.log(`[glacier-poller] failed to process restore ${row.uuid} (${row.videoId}): ${err.message}`) + await db.updateTable('restore_requests') + .set({ status: 'failed', error_message: err.message, updated_at: new Date() }) + .where('uuid', '=', row.uuid) + .execute() + } + } +} + +function startRestorePoller() { + pollRestores() + setInterval(pollRestores, POLL_INTERVAL_MS).unref() + console.log('glacier restore poller started (this server is primary)') +} + +export { startRestorePoller, pollRestores } diff --git a/src/utils/mail.ts b/src/utils/mail.ts new file mode 100644 index 0000000..b8a0e51 --- /dev/null +++ b/src/utils/mail.ts @@ -0,0 +1,23 @@ +import nodemailer from 'nodemailer' + +const transporter = nodemailer.createTransport({ + host: process.env.SMTP_HOST, + port: Number(process.env.SMTP_PORT || 587), + auth: { + user: process.env.SMTP_USER, + pass: process.env.SMTP_PASS + } +}) + +async function sendRestoreCompleteEmail(to: string, videoId: string, title?: string | null) { + const watchUrl = `https://preservetube.com/watch?v=${videoId}` + + await transporter.sendMail({ + from: process.env.SMTP_FROM, + to, + subject: 'Your PreserveTube video has been restored', + text: `The video you requested${title ? ` ("${title}")` : ''} has been restored from cold storage and is available again:\n\n${watchUrl}` + }) +} + +export { sendRestoreCompleteEmail }