mirror of
https://github.com/bitinflow/rerun-encoder.git
synced 2026-03-13 13:46:00 +00:00
Rewrite s3 uploader
Update user on upload/boot
This commit is contained in:
@@ -1,8 +1,27 @@
|
||||
import {BrowserWindow} from "electron";
|
||||
import {User} from "../../shared/schema";
|
||||
import axios from "axios";
|
||||
|
||||
export function emit(event: any, ...args: any) {
|
||||
// Send a message to all windows
|
||||
BrowserWindow.getAllWindows().forEach((win) => {
|
||||
win.webContents.send(event, ...args)
|
||||
});
|
||||
}
|
||||
|
||||
export async function resolveUser(accessToken: string, tokenType: string): Promise<User> {
|
||||
const response = await axios.get('https://api.rerunmanager.com/v1/channels/me', {
|
||||
headers: {
|
||||
'Accept': 'application/json',
|
||||
'Authorization': `${tokenType} ${accessToken}`,
|
||||
}
|
||||
});
|
||||
|
||||
return {
|
||||
id: response.data.id,
|
||||
name: response.data.name,
|
||||
config: response.data.config,
|
||||
avatar_url: response.data.avatar_url,
|
||||
premium: response.data.premium,
|
||||
}
|
||||
}
|
||||
@@ -161,8 +161,8 @@ ipcMain.handle('encode', async (event: IpcMainInvokeEvent, ...args: any[]) => {
|
||||
onUploadProgress: (id, progress) => event.sender.send('encode-upload-progress', id, progress),
|
||||
onUploadComplete: (id, video) => event.sender.send('encode-upload-complete', id, video),
|
||||
onError: (id, error) => event.sender.send('encode-error', id, error),
|
||||
}, settingsRepository.getSettings());
|
||||
return await encoder.encode()
|
||||
}, settingsRepository);
|
||||
return encoder.encode();
|
||||
})
|
||||
ipcMain.handle('commitSettings', async (event: IpcMainInvokeEvent, ...args: any[]) => settingsRepository.commitSettings(args[0]))
|
||||
settingsRepository.watch((settings: Settings) => emit('settings', settings));
|
||||
|
||||
@@ -1,7 +1,10 @@
|
||||
|
||||
import {EncoderListeners, EncoderOptions, Settings, User, Video} from "../../shared/schema";
|
||||
import {EncoderListeners, EncoderOptions, Video} from "../../shared/schema";
|
||||
import * as fs from "fs";
|
||||
import axios, {AxiosInstance} from "axios";
|
||||
import {SettingsRepository} from "./settings-repository";
|
||||
import * as path from "path";
|
||||
import {path as ffmpegPath} from "@ffmpeg-installer/ffmpeg";
|
||||
import {upload} from "./file-uploader";
|
||||
|
||||
export class Encoder {
|
||||
private readonly id: string;
|
||||
@@ -9,7 +12,7 @@ export class Encoder {
|
||||
private readonly output: string;
|
||||
private readonly options: EncoderOptions;
|
||||
private readonly listeners: EncoderListeners;
|
||||
private readonly settings: Settings;
|
||||
private readonly settingsRepository: SettingsRepository;
|
||||
private api: AxiosInstance;
|
||||
private s3: AxiosInstance;
|
||||
|
||||
@@ -18,15 +21,17 @@ export class Encoder {
|
||||
input: string,
|
||||
options: EncoderOptions,
|
||||
listeners: EncoderListeners,
|
||||
settings: Settings
|
||||
settingsRepository: SettingsRepository
|
||||
) {
|
||||
|
||||
this.id = id;
|
||||
this.input = input;
|
||||
this.output = this.input.replace(/\.mp4$/, '.flv')
|
||||
this.output = this.getOutputPath(input, 'flv');
|
||||
this.options = options;
|
||||
this.listeners = listeners;
|
||||
this.settings = settings;
|
||||
this.settingsRepository = settingsRepository;
|
||||
|
||||
const settings = settingsRepository.getSettings();
|
||||
this.api = axios.create({
|
||||
baseURL: settings.endpoint,
|
||||
headers: {
|
||||
@@ -37,45 +42,36 @@ export class Encoder {
|
||||
this.s3 = axios.create({});
|
||||
}
|
||||
|
||||
async encode(): Promise<void> {
|
||||
encode(): void {
|
||||
this.listeners.onStart(this.id)
|
||||
|
||||
const ffmpegPath = require('@ffmpeg-installer/ffmpeg').path
|
||||
console.log('ffmpegPath', ffmpegPath)
|
||||
const ffmpeg = require('fluent-ffmpeg')
|
||||
ffmpeg.setFfmpegPath(ffmpegPath)
|
||||
let totalTime = 0;
|
||||
ffmpeg(this.input)
|
||||
.outputOptions(this.getOutputOptions())
|
||||
.output(this.output)
|
||||
.on('start', () => {
|
||||
console.log('start')
|
||||
})
|
||||
.on('codecData', data => {
|
||||
totalTime = parseInt(data.duration.replace(/:/g, ''))
|
||||
})
|
||||
.on('progress', progress => {
|
||||
const time = parseInt(progress.timemark.replace(/:/g, ''))
|
||||
const percent = (time / totalTime) * 100
|
||||
console.log('progress', percent)
|
||||
this.listeners.onProgress(this.id, percent)
|
||||
})
|
||||
.on('end', async () => {
|
||||
console.log('end')
|
||||
// @ts-ignore
|
||||
try {
|
||||
const video = await this.requestUploadUrl(this.settings.credentials.user)
|
||||
await this.upload(video)
|
||||
} catch (error) {
|
||||
if (this.requireEncoding()) {
|
||||
const ffmpegPath = require('@ffmpeg-installer/ffmpeg').path
|
||||
console.log('ffmpegPath', ffmpegPath)
|
||||
const ffmpeg = require('fluent-ffmpeg')
|
||||
ffmpeg.setFfmpegPath(ffmpegPath)
|
||||
let totalTime = 0;
|
||||
|
||||
ffmpeg(this.input)
|
||||
.output(this.output)
|
||||
.outputOptions(this.getOutputOptions())
|
||||
.on('start', () => console.log('start'))
|
||||
.on('codecData', data => totalTime = parseInt(data.duration.replace(/:/g, '')))
|
||||
.on('progress', progress => {
|
||||
const time = parseInt(progress.timemark.replace(/:/g, ''))
|
||||
const percent = (time / totalTime) * 100
|
||||
console.log('progress', percent)
|
||||
this.listeners.onProgress(this.id, percent)
|
||||
})
|
||||
.on('end', async () => this.requestUpload(this.output))
|
||||
.on('error', (error) => {
|
||||
console.log('error', error)
|
||||
this.listeners.onError(this.id, error.message)
|
||||
}
|
||||
})
|
||||
.on('error', (error) => {
|
||||
console.log('error', error)
|
||||
this.listeners.onError(this.id, error.message)
|
||||
})
|
||||
.run()
|
||||
})
|
||||
.run()
|
||||
} else {
|
||||
this.requestUpload(this.input)
|
||||
}
|
||||
}
|
||||
|
||||
private getOutputOptions() {
|
||||
@@ -94,13 +90,21 @@ export class Encoder {
|
||||
]
|
||||
}
|
||||
|
||||
private async requestUploadUrl(user: User): Promise<Video> {
|
||||
private async requestUpload(filename: string): Promise<void> {
|
||||
const user = this.settingsRepository.getSettings().credentials.user;
|
||||
const response = await this.api.post(`channels/${user.id}/videos`, {
|
||||
title: this.options.title,
|
||||
size: fs.statSync(this.output).size,
|
||||
size: fs.statSync(filename).size,
|
||||
})
|
||||
|
||||
return response.data as Video
|
||||
const video = response.data as Video
|
||||
|
||||
try {
|
||||
await this.handleUpload(video, filename)
|
||||
} catch (error) {
|
||||
console.log('error', error)
|
||||
this.listeners.onError(this.id, error.message)
|
||||
}
|
||||
}
|
||||
|
||||
private async confirmUpload(video: Video): Promise<Video> {
|
||||
@@ -117,30 +121,18 @@ export class Encoder {
|
||||
await this.api.delete(`videos/${video.id}/cancel`)
|
||||
}
|
||||
|
||||
private async upload(video: Video) {
|
||||
private async handleUpload(video: Video, filename: string) {
|
||||
if (!video.upload_url) {
|
||||
return
|
||||
}
|
||||
|
||||
const stream = fs.createReadStream(this.output)
|
||||
stream.on('error', (error) => {
|
||||
console.log('upload error', error)
|
||||
this.listeners.onError(this.id, error.message)
|
||||
})
|
||||
|
||||
const progress = (progress: any) => {
|
||||
const progressCompleted = progress.loaded / fs.statSync(this.output).size * 100
|
||||
this.listeners.onUploadProgress(this.id, progressCompleted)
|
||||
}
|
||||
|
||||
try {
|
||||
await this.s3.put(video.upload_url, stream, {
|
||||
onUploadProgress: progress,
|
||||
headers: {
|
||||
'Content-Type': 'video/x-flv',
|
||||
}
|
||||
await upload(video.upload_url, filename, (progress: number) => {
|
||||
this.listeners.onUploadProgress(this.id, progress)
|
||||
})
|
||||
|
||||
this.settingsRepository.reloadUser();
|
||||
|
||||
this.listeners.onUploadComplete(this.id, await this.confirmUpload(video))
|
||||
} catch (error) {
|
||||
console.log('upload error', error)
|
||||
@@ -152,4 +144,22 @@ export class Encoder {
|
||||
this.listeners.onError(this.id, `Upload Error: ${error.message}`)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* This will return the output path for the encoded file.
|
||||
* It will have the same name as the input file, but with the given extension.
|
||||
* Also, a suffix will be added to the file name to avoid overwriting existing files.
|
||||
*/
|
||||
private getOutputPath(input: string, ext: string) {
|
||||
const parsed = path.parse(input)
|
||||
return path.join(parsed.dir, `${parsed.name}-rerun.${ext}`)
|
||||
}
|
||||
|
||||
/**
|
||||
* You can also just upload the file as is, without encoding it.
|
||||
* @private
|
||||
*/
|
||||
private requireEncoding() {
|
||||
return path.parse(this.input).ext !== '.flv'
|
||||
}
|
||||
}
|
||||
39
electron/rerun-manager/file-uploader.ts
Normal file
39
electron/rerun-manager/file-uploader.ts
Normal file
@@ -0,0 +1,39 @@
|
||||
const fs = require('fs');
|
||||
const https = require('https');
|
||||
const {promisify} = require('util');
|
||||
|
||||
export async function upload(url, filename, onProgress): Promise<void> {
|
||||
const fileStream = fs.createReadStream(filename);
|
||||
const fileStats = await promisify(fs.stat)(filename);
|
||||
|
||||
return new Promise((resolve, reject) => {
|
||||
const options = {
|
||||
method: 'PUT',
|
||||
headers: {
|
||||
'Content-Type': 'video/x-flv',
|
||||
'Content-Length': fileStats.size,
|
||||
},
|
||||
agent: false // Disable HTTP keep-alive
|
||||
}
|
||||
|
||||
const req = https.request(url, options, (res) => {
|
||||
if (res.statusCode >= 400) {
|
||||
reject(new Error(`Failed to upload file: ${res.statusCode} ${res.statusMessage}`));
|
||||
return;
|
||||
}
|
||||
|
||||
resolve();
|
||||
});
|
||||
|
||||
req.on('error', reject);
|
||||
|
||||
let uploadedBytes = 0;
|
||||
fileStream.on('data', (chunk) => {
|
||||
uploadedBytes += chunk.length;
|
||||
onProgress(Math.round((uploadedBytes / fileStats.size) * 100));
|
||||
});
|
||||
|
||||
fileStream.on('error', reject);
|
||||
fileStream.pipe(req);
|
||||
});
|
||||
}
|
||||
@@ -6,10 +6,11 @@ import koaCors from "koa-cors";
|
||||
import bodyParser from "koa-bodyparser";
|
||||
import axios, {AxiosInstance} from "axios";
|
||||
import {app} from "electron"
|
||||
import {Credentials, User} from "../../shared/schema";
|
||||
import {Credentials} from "../../shared/schema";
|
||||
import {SettingsRepository} from "./settings-repository";
|
||||
import * as fs from "fs";
|
||||
import {join} from "node:path";
|
||||
import {resolveUser} from "../main/helpers";
|
||||
|
||||
export class InternalServer {
|
||||
private readonly app: Application;
|
||||
@@ -50,7 +51,7 @@ export class InternalServer {
|
||||
expires_in: params.get('expires_in'),
|
||||
expires_at: this.calculateExpiresAt(params.get('expires_in')),
|
||||
state: params.get('state'),
|
||||
user: await this.resolveUser(params.get('access_token'), params.get('token_type')),
|
||||
user: await resolveUser(params.get('access_token'), params.get('token_type')),
|
||||
}
|
||||
|
||||
console.log('credentials', credentials);
|
||||
@@ -76,23 +77,6 @@ export class InternalServer {
|
||||
});
|
||||
}
|
||||
|
||||
private async resolveUser(accessToken: string, tokenType: string): Promise<User> {
|
||||
const response = await this.axios.get('https://api.rerunmanager.com/v1/channels/me', {
|
||||
headers: {
|
||||
'Accept': 'application/json',
|
||||
'Authorization': `${tokenType} ${accessToken}`,
|
||||
}
|
||||
});
|
||||
|
||||
return {
|
||||
id: response.data.id,
|
||||
name: response.data.name,
|
||||
config: response.data.config,
|
||||
avatar_url: response.data.avatar_url,
|
||||
premium: response.data.premium,
|
||||
}
|
||||
}
|
||||
|
||||
private calculateExpiresAt(expiresIn: string) {
|
||||
const now = new Date();
|
||||
const expiresAt = new Date(now.getTime() + parseInt(expiresIn) * 1000);
|
||||
|
||||
@@ -2,6 +2,7 @@ import defu from 'defu'
|
||||
import {platform} from 'node:process'
|
||||
import * as fs from 'fs'
|
||||
import {Credentials, Settings} from "../../shared/schema";
|
||||
import {resolveUser} from "../main/helpers";
|
||||
|
||||
const defaults: Settings = {
|
||||
version: '1.0.1',
|
||||
@@ -59,6 +60,8 @@ export class SettingsRepository {
|
||||
// if so, set settings.credentials to null
|
||||
if (this.isExpired()) {
|
||||
console.log('Credentials expired!');
|
||||
} else {
|
||||
this.reloadUser();
|
||||
}
|
||||
|
||||
await this.save()
|
||||
@@ -121,4 +124,19 @@ export class SettingsRepository {
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
reloadUser() {
|
||||
try {
|
||||
console.debug('Reloading user')
|
||||
resolveUser(
|
||||
this.settings.credentials.access_token,
|
||||
this.settings.credentials.token_type
|
||||
).then((user) => {
|
||||
this.settings.credentials.user = user
|
||||
this.save()
|
||||
})
|
||||
} catch (e) {
|
||||
console.error('Failed to reload user', e)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user