211 lines
7.4 KiB
Go
211 lines
7.4 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
|
)
|
|
|
|
const maxCreatorMaterialBytes int64 = 512 << 20
|
|
|
|
func processCreatorMaterial(ctx context.Context, store *creator.Store, workID string) (creator.MaterialJob, error) {
|
|
if store == nil || workID == "" || filepath.Base(workID) != workID {
|
|
return creator.MaterialJob{}, creator.ErrInvalid
|
|
}
|
|
work, err := store.GetWork(ctx, workID)
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
job, err := store.GetMaterial(ctx, workID)
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
if !job.Selected {
|
|
return creator.MaterialJob{}, creator.ErrConflict
|
|
}
|
|
root := os.Getenv("CREATOR_MEDIA_DIR")
|
|
if root == "" {
|
|
root = "/var/lib/creatorhub/materials"
|
|
}
|
|
root = filepath.Clean(root)
|
|
dir := filepath.Join(root, workID)
|
|
if err := os.MkdirAll(dir, 0o700); err != nil {
|
|
return creator.MaterialJob{}, fmt.Errorf("create material directory: %w", err)
|
|
}
|
|
videoPath := filepath.Join(dir, "source")
|
|
videoReference := filepath.Join(workID, "source")
|
|
if job.DownloadStatus != "succeeded" || !fileExists(videoPath) {
|
|
if _, err := store.SetMaterialStep(ctx, workID, "download", "running", "", ""); err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
if err := downloadCreatorMaterial(ctx, work.OriginalURL, videoPath); err != nil {
|
|
return setMaterialFailure(ctx, store, workID, "download", err)
|
|
}
|
|
job, err = store.SetMaterialStep(ctx, workID, "download", "succeeded", videoReference, "")
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
}
|
|
|
|
audioPath := filepath.Join(dir, "audio.wav")
|
|
audioReference := filepath.Join(workID, "audio.wav")
|
|
if job.AudioStatus != "succeeded" && job.AudioStatus != "no_audio" {
|
|
if _, err := store.SetMaterialStep(ctx, workID, "audio", "running", "", ""); err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
hasAudio, err := creatorMaterialHasAudio(ctx, videoPath)
|
|
if err != nil {
|
|
return setMaterialFailure(ctx, store, workID, "audio", err)
|
|
}
|
|
if !hasAudio {
|
|
job, err = store.SetMaterialStep(ctx, workID, "audio", "no_audio", "", "视频没有音轨")
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
} else if err := extractCreatorAudio(ctx, videoPath, audioPath); err != nil {
|
|
return setMaterialFailure(ctx, store, workID, "audio", err)
|
|
} else {
|
|
job, err = store.SetMaterialStep(ctx, workID, "audio", "succeeded", audioReference, "")
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
}
|
|
}
|
|
|
|
if job.TranscriptionStatus != "succeeded" && job.TranscriptionStatus != "no_speech" {
|
|
if job.AudioStatus == "no_audio" {
|
|
job, err = store.SetMaterialStep(ctx, workID, "transcription", "no_speech", "", "没有可转写的音轨")
|
|
} else {
|
|
if _, err = store.SetMaterialStep(ctx, workID, "transcription", "running", "", ""); err == nil {
|
|
transcript, transcribeErr := transcribeCreatorAudio(ctx, audioPath)
|
|
if transcribeErr != nil {
|
|
job, err = setMaterialFailure(ctx, store, workID, "transcription", transcribeErr)
|
|
} else if strings.TrimSpace(transcript) == "" {
|
|
job, err = store.SetMaterialStep(ctx, workID, "transcription", "no_speech", "", "转写未检测到语音")
|
|
} else {
|
|
transcriptPath := filepath.Join(dir, "transcript.txt")
|
|
if writeErr := os.WriteFile(transcriptPath, []byte(transcript), 0o600); writeErr != nil {
|
|
job, err = setMaterialFailure(ctx, store, workID, "transcription", writeErr)
|
|
} else {
|
|
job, err = store.SetMaterialStep(ctx, workID, "transcription", "succeeded", filepath.Join(workID, "transcript.txt"), "")
|
|
}
|
|
}
|
|
}
|
|
}
|
|
if err != nil {
|
|
return creator.MaterialJob{}, err
|
|
}
|
|
}
|
|
return job, nil
|
|
}
|
|
|
|
func downloadCreatorMaterial(ctx context.Context, rawURL, destination string) error {
|
|
parsed, err := http.NewRequestWithContext(ctx, http.MethodGet, rawURL, nil)
|
|
if err != nil || (parsed.URL.Scheme != "http" && parsed.URL.Scheme != "https") || parsed.URL.Host == "" {
|
|
return creator.ErrInvalid
|
|
}
|
|
client := &http.Client{Timeout: 2 * time.Minute}
|
|
response, err := client.Do(parsed)
|
|
if err != nil {
|
|
return fmt.Errorf("download material: %w", err)
|
|
}
|
|
defer response.Body.Close()
|
|
if response.StatusCode < 200 || response.StatusCode >= 300 {
|
|
return fmt.Errorf("download material: HTTP %s", response.Status)
|
|
}
|
|
if response.ContentLength > maxCreatorMaterialBytes {
|
|
return fmt.Errorf("download material exceeds size limit")
|
|
}
|
|
temporary := destination + ".tmp"
|
|
defer os.Remove(temporary)
|
|
file, err := os.OpenFile(temporary, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600)
|
|
if err != nil {
|
|
return fmt.Errorf("create material file: %w", err)
|
|
}
|
|
written, copyErr := io.Copy(file, io.LimitReader(response.Body, maxCreatorMaterialBytes+1))
|
|
closeErr := file.Close()
|
|
if copyErr != nil {
|
|
return fmt.Errorf("write material: %w", copyErr)
|
|
}
|
|
if closeErr != nil {
|
|
return fmt.Errorf("close material file: %w", closeErr)
|
|
}
|
|
if written > maxCreatorMaterialBytes {
|
|
return fmt.Errorf("download material exceeds size limit")
|
|
}
|
|
if err := os.Rename(temporary, destination); err != nil {
|
|
return fmt.Errorf("publish material: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func creatorMaterialHasAudio(ctx context.Context, videoPath string) (bool, error) {
|
|
if _, err := exec.LookPath("ffprobe"); err != nil {
|
|
return false, fmt.Errorf("ffprobe is unavailable: %w", err)
|
|
}
|
|
command := exec.CommandContext(ctx, "ffprobe", "-v", "error", "-select_streams", "a:0", "-show_entries", "stream=index", "-of", "csv=p=0", videoPath)
|
|
output, err := command.Output()
|
|
if err != nil {
|
|
return false, fmt.Errorf("inspect audio stream: %w", err)
|
|
}
|
|
return strings.TrimSpace(string(output)) != "", nil
|
|
}
|
|
|
|
func extractCreatorAudio(ctx context.Context, videoPath, audioPath string) error {
|
|
if _, err := exec.LookPath("ffmpeg"); err != nil {
|
|
return fmt.Errorf("ffmpeg is unavailable: %w", err)
|
|
}
|
|
temporary := audioPath + ".tmp"
|
|
defer os.Remove(temporary)
|
|
command := exec.CommandContext(ctx, "ffmpeg", "-nostdin", "-v", "error", "-y", "-i", videoPath, "-vn", "-ac", "1", "-ar", "16000", temporary)
|
|
if output, err := command.CombinedOutput(); err != nil {
|
|
return fmt.Errorf("extract audio: %w: %s", err, strings.TrimSpace(string(output)))
|
|
}
|
|
if !fileExists(temporary) {
|
|
return fmt.Errorf("extract audio produced no file")
|
|
}
|
|
if err := os.Rename(temporary, audioPath); err != nil {
|
|
return fmt.Errorf("publish audio: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func transcribeCreatorAudio(ctx context.Context, audioPath string) (string, error) {
|
|
binary := os.Getenv("CREATOR_TRANSCRIPTION_BIN")
|
|
if binary == "" {
|
|
return "", fmt.Errorf("transcription provider is not configured")
|
|
}
|
|
if _, err := exec.LookPath(binary); err != nil {
|
|
return "", fmt.Errorf("transcription provider is unavailable: %w", err)
|
|
}
|
|
output, err := exec.CommandContext(ctx, binary, audioPath).Output()
|
|
if err != nil {
|
|
return "", fmt.Errorf("transcribe audio: %w", err)
|
|
}
|
|
if len(output) > 1<<20 {
|
|
return "", fmt.Errorf("transcript exceeds size limit")
|
|
}
|
|
return string(output), nil
|
|
}
|
|
|
|
func setMaterialFailure(ctx context.Context, store *creator.Store, workID, step string, cause error) (creator.MaterialJob, error) {
|
|
job, err := store.SetMaterialStep(ctx, workID, step, "failed", "", cause.Error())
|
|
if err != nil {
|
|
return creator.MaterialJob{}, fmt.Errorf("record %s failure: %w", step, err)
|
|
}
|
|
return job, nil
|
|
}
|
|
|
|
func fileExists(path string) bool {
|
|
info, err := os.Stat(path)
|
|
return err == nil && info.Mode().IsRegular() && info.Size() > 0
|
|
}
|