feat: Add sequential downloading and streaming support via HTTP server
This commit is contained in:
parent
36972d61b6
commit
14860d28b2
5 changed files with 278 additions and 45 deletions
|
|
@ -6,6 +6,7 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/veggiedefender/torrent-client/internal/history"
|
||||
"github.com/veggiedefender/torrent-client/internal/stream"
|
||||
"github.com/veggiedefender/torrent-client/internal/torrent"
|
||||
)
|
||||
|
||||
|
|
@ -15,6 +16,8 @@ type Controller struct {
|
|||
cancel context.CancelFunc
|
||||
currentPath string
|
||||
currentOut string
|
||||
|
||||
streamServer *stream.Server
|
||||
}
|
||||
|
||||
func NewController() *Controller {
|
||||
|
|
@ -25,6 +28,7 @@ func NewController() *Controller {
|
|||
engine: engine,
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
streamServer: stream.NewServer(engine),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -98,3 +102,15 @@ func (c *Controller) Progress() float64 {
|
|||
func (c *Controller) Status() torrent.Status {
|
||||
return c.engine.Status()
|
||||
}
|
||||
|
||||
func (c *Controller) SetSequentialMode(mode bool) {
|
||||
c.engine.SetSequentialMode(mode)
|
||||
}
|
||||
|
||||
func (c *Controller) StartStreamServer() {
|
||||
c.streamServer.Start()
|
||||
}
|
||||
|
||||
func (c *Controller) StopStreamServer() {
|
||||
c.streamServer.Stop()
|
||||
}
|
||||
|
|
|
|||
142
internal/stream/server.go
Normal file
142
internal/stream/server.go
Normal file
|
|
@ -0,0 +1,142 @@
|
|||
package stream
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/veggiedefender/torrent-client/internal/torrent"
|
||||
)
|
||||
|
||||
// Server handles HTTP streaming for a torrent
|
||||
type Server struct {
|
||||
server *http.Server
|
||||
engine *torrent.Engine
|
||||
}
|
||||
|
||||
func NewServer(engine *torrent.Engine) *Server {
|
||||
s := &Server{
|
||||
engine: engine,
|
||||
}
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/stream", s.handleStream)
|
||||
|
||||
s.server = &http.Server{
|
||||
Addr: ":8080",
|
||||
Handler: mux,
|
||||
}
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
func (s *Server) Start() {
|
||||
go func() {
|
||||
log.Println("Starting stream server on http://localhost:8080/stream")
|
||||
if err := s.server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
log.Printf("Stream server error: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (s *Server) Stop() {
|
||||
if s.server != nil {
|
||||
s.server.Shutdown(context.Background())
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) handleStream(w http.ResponseWriter, r *http.Request) {
|
||||
status := s.engine.Status()
|
||||
if !status.Loaded {
|
||||
http.Error(w, "Torrent not loaded", http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
|
||||
partPath := s.engine.PartFilePath()
|
||||
if partPath == "" {
|
||||
http.Error(w, "Part file not ready", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
|
||||
file, err := os.Open(partPath)
|
||||
if err != nil {
|
||||
http.Error(w, "Failed to open file", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// Note: We don't defer file.Close() here because ServeContent uses the Seeker asynchronously,
|
||||
// wait, http.ServeContent does NOT close it, but it finishes before returning!
|
||||
// So we MUST defer file.Close().
|
||||
defer file.Close()
|
||||
|
||||
reader := &streamReader{
|
||||
file: file,
|
||||
engine: s.engine,
|
||||
pieceLength: status.PieceLength,
|
||||
length: int64(status.Length),
|
||||
offset: 0,
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", "video/mp4") // TODO: detect properly, default to mp4
|
||||
w.Header().Set("Accept-Ranges", "bytes")
|
||||
|
||||
http.ServeContent(w, r, status.Name, time.Now(), reader)
|
||||
}
|
||||
|
||||
// streamReader wraps a file and blocks until requested pieces are downloaded
|
||||
type streamReader struct {
|
||||
file *os.File
|
||||
engine *torrent.Engine
|
||||
pieceLength int
|
||||
length int64
|
||||
offset int64
|
||||
}
|
||||
|
||||
func (r *streamReader) Read(p []byte) (n int, err error) {
|
||||
if r.offset >= r.length {
|
||||
return 0, io.EOF
|
||||
}
|
||||
|
||||
// Calculate which piece we are trying to read
|
||||
pieceIndex := int(r.offset / int64(r.pieceLength))
|
||||
|
||||
// Polling loop to wait for the piece
|
||||
for !r.engine.HasPiece(pieceIndex) {
|
||||
status := r.engine.Status()
|
||||
if status.Phase == "failed" || status.Phase == "stopped" {
|
||||
return 0, fmt.Errorf("torrent stopped or failed")
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
|
||||
// Read from the actual file
|
||||
n, err = r.file.ReadAt(p, r.offset)
|
||||
if n > 0 {
|
||||
r.offset += int64(n)
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
func (r *streamReader) Seek(offset int64, whence int) (int64, error) {
|
||||
var newOffset int64
|
||||
switch whence {
|
||||
case io.SeekStart:
|
||||
newOffset = offset
|
||||
case io.SeekCurrent:
|
||||
newOffset = r.offset + offset
|
||||
case io.SeekEnd:
|
||||
newOffset = r.length + offset
|
||||
default:
|
||||
return 0, fmt.Errorf("invalid whence: %d", whence)
|
||||
}
|
||||
|
||||
if newOffset < 0 {
|
||||
return 0, fmt.Errorf("negative offset")
|
||||
}
|
||||
|
||||
r.offset = newOffset
|
||||
return r.offset, nil
|
||||
}
|
||||
|
|
@ -66,6 +66,8 @@ type Engine struct {
|
|||
// rate limiting
|
||||
downloadLimitBps atomic.Int64
|
||||
uploadLimitBps atomic.Int64
|
||||
|
||||
scheduler *pieceScheduler // For toggling sequential mode
|
||||
}
|
||||
|
||||
type PieceState uint8
|
||||
|
|
@ -736,6 +738,9 @@ func (e *Engine) runDownload(ctx context.Context, tf *torrentfile.TorrentFile) {
|
|||
e.setTerminalError("failed", fmt.Errorf("create temp file: %w", err))
|
||||
return
|
||||
}
|
||||
e.mu.Lock()
|
||||
e.outputPath = partPath
|
||||
e.mu.Unlock()
|
||||
defer partFile.Close()
|
||||
|
||||
// Восстанавливаем состояние уже скачанных кусков
|
||||
|
|
@ -759,6 +764,15 @@ func (e *Engine) runDownload(ctx context.Context, tf *torrentfile.TorrentFile) {
|
|||
defer pieceWriter.Close()
|
||||
|
||||
scheduler := newPieceSchedulerWithResume(len(tf.PieceHashes), resumedPieces)
|
||||
e.mu.Lock()
|
||||
e.scheduler = scheduler
|
||||
e.mu.Unlock()
|
||||
defer func() {
|
||||
e.mu.Lock()
|
||||
e.scheduler = nil
|
||||
e.mu.Unlock()
|
||||
scheduler.Stop()
|
||||
}()
|
||||
workerCtx, workerCancel := context.WithCancel(ctx)
|
||||
var workersWG sync.WaitGroup
|
||||
|
||||
|
|
@ -2243,3 +2257,29 @@ func (e *Engine) handleIncomingSeeding(ctx context.Context, tf *torrentfile.Torr
|
|||
}()
|
||||
}
|
||||
|
||||
// SetSequentialMode toggles sequential piece downloading mode on the fly
|
||||
func (e *Engine) SetSequentialMode(mode bool) {
|
||||
e.mu.RLock()
|
||||
scheduler := e.scheduler
|
||||
e.mu.RUnlock()
|
||||
if scheduler != nil {
|
||||
scheduler.SetSequential(mode)
|
||||
}
|
||||
}
|
||||
|
||||
// HasPiece is a thread-safe check to see if a piece is fully downloaded
|
||||
func (e *Engine) HasPiece(index int) bool {
|
||||
e.mu.RLock()
|
||||
defer e.mu.RUnlock()
|
||||
if index < 0 || index >= len(e.pieceStates) {
|
||||
return false
|
||||
}
|
||||
return e.pieceStates[index] == PieceCompleted
|
||||
}
|
||||
|
||||
// PartFilePath returns the path to the current .part file
|
||||
func (e *Engine) PartFilePath() string {
|
||||
e.mu.RLock()
|
||||
defer e.mu.RUnlock()
|
||||
return e.outputPath
|
||||
}
|
||||
|
|
|
|||
|
|
@ -42,6 +42,8 @@ type pieceScheduler struct {
|
|||
stopCh chan struct{}
|
||||
stopOnce sync.Once
|
||||
|
||||
setSequentialCh chan bool
|
||||
|
||||
// cancelCh рассылает сигнал отмены конкретному воркеру в endgame-режиме.
|
||||
// Ключ — peerID, значение — канал с индексом куска для отмены.
|
||||
cancelMu sync.RWMutex
|
||||
|
|
@ -68,9 +70,10 @@ func newPieceSchedulerWithResume(pieceCount int, completedIndices []int) *pieceS
|
|||
assignCh: make(chan assignPieceRequest, 128),
|
||||
reportCh: make(chan reportPieceRequest, 128),
|
||||
releasePeerCh: make(chan string, 128),
|
||||
progressCh: make(chan int, 128),
|
||||
progressCh: make(chan int, 1),
|
||||
doneCh: make(chan struct{}),
|
||||
stopCh: make(chan struct{}),
|
||||
setSequentialCh: make(chan bool),
|
||||
cancelSubs: make(map[string]chan int),
|
||||
}
|
||||
|
||||
|
|
@ -189,15 +192,21 @@ func (ps *pieceScheduler) Stop() {
|
|||
})
|
||||
}
|
||||
|
||||
func (ps *pieceScheduler) SetSequential(mode bool) {
|
||||
select {
|
||||
case ps.setSequentialCh <- mode:
|
||||
case <-ps.stopCh:
|
||||
}
|
||||
}
|
||||
|
||||
func (ps *pieceScheduler) run(pieceCount int, completedIndices []int) {
|
||||
states := make([]pieceState, pieceCount)
|
||||
availability := make([]int, pieceCount)
|
||||
peerAvailability := make(map[string][]bool, 256)
|
||||
// endgamePeers — множество пиров, которым назначен кусок в endgame режиме
|
||||
endgamePeers := make(map[int][]string) // pieceIndex -> []peerID
|
||||
endgamePeers := make(map[int][]string)
|
||||
completed := 0
|
||||
sequentialMode := false
|
||||
|
||||
// Восстанавливаем уже завершённые куски из resume
|
||||
for _, idx := range completedIndices {
|
||||
if idx >= 0 && idx < pieceCount {
|
||||
states[idx] = pieceDone
|
||||
|
|
@ -221,7 +230,6 @@ func (ps *pieceScheduler) run(pieceCount int, completedIndices []int) {
|
|||
return
|
||||
}
|
||||
|
||||
// Подсчёт оставшихся кусков для определения endgame
|
||||
remaining := pieceCount - completed
|
||||
endgame := remaining <= endgameThreshold
|
||||
|
||||
|
|
@ -238,19 +246,20 @@ func (ps *pieceScheduler) run(pieceCount int, completedIndices []int) {
|
|||
|
||||
pieceIndex := -1
|
||||
if endgame {
|
||||
// Endgame: разрешаем брать куски в состоянии InProgress тоже
|
||||
pieceIndex = selectPieceEndgame(states, availability, req.have, req.hasInfo, endgamePeers, req.peerID)
|
||||
} else {
|
||||
if sequentialMode {
|
||||
pieceIndex = selectPendingPieceSequential(states, req.have, req.hasInfo)
|
||||
} else {
|
||||
pieceIndex = selectPendingPieceRarest(states, availability, req.have, req.hasInfo)
|
||||
}
|
||||
}
|
||||
|
||||
if pieceIndex >= 0 {
|
||||
if !endgame {
|
||||
states[pieceIndex] = pieceInProgress
|
||||
}
|
||||
// В endgame: запоминаем всех пиров, которые качают этот кусок
|
||||
endgamePeers[pieceIndex] = append(endgamePeers[pieceIndex], req.peerID)
|
||||
debugf("scheduler assigned piece %d to %s (endgame=%v)", pieceIndex, req.peerID, endgame)
|
||||
req.responseCh <- assignPieceResponse{
|
||||
task: pieceTask{Index: pieceIndex},
|
||||
ok: true,
|
||||
|
|
@ -260,25 +269,21 @@ func (ps *pieceScheduler) run(pieceCount int, completedIndices []int) {
|
|||
req.responseCh <- assignPieceResponse{ok: false}
|
||||
|
||||
case req := <-ps.reportCh:
|
||||
if req.pieceIndex < 0 || req.pieceIndex >= len(states) {
|
||||
if req.pieceIndex < 0 || req.pieceIndex >= pieceCount {
|
||||
req.responseCh <- false
|
||||
continue
|
||||
}
|
||||
|
||||
if req.success {
|
||||
if states[req.pieceIndex] == pieceDone {
|
||||
debugf("scheduler report piece %d ignored (already done)", req.pieceIndex)
|
||||
req.responseCh <- false
|
||||
continue
|
||||
}
|
||||
states[req.pieceIndex] = pieceDone
|
||||
completed++
|
||||
debugf("scheduler report piece %d success (%d/%d)", req.pieceIndex, completed, pieceCount)
|
||||
|
||||
// Endgame: уведомить других пиров отменить этот кусок
|
||||
if peers, ok := endgamePeers[req.pieceIndex]; ok && len(peers) > 1 {
|
||||
// Находим winner — последний репортнувший (req не содержит peerID,
|
||||
// поэтому рассылаем всем — воркер проверяет индекс)
|
||||
go ps.broadcastCancel(req.pieceIndex, "")
|
||||
}
|
||||
delete(endgamePeers, req.pieceIndex)
|
||||
|
|
@ -287,30 +292,15 @@ func (ps *pieceScheduler) run(pieceCount int, completedIndices []int) {
|
|||
case ps.progressCh <- req.pieceIndex:
|
||||
default:
|
||||
}
|
||||
req.responseCh <- true
|
||||
continue
|
||||
}
|
||||
|
||||
// Неуспех: возвращаем кусок в pending
|
||||
if states[req.pieceIndex] == pieceInProgress {
|
||||
states[req.pieceIndex] = piecePending
|
||||
debugf("scheduler report piece %d failed, re-queued", req.pieceIndex)
|
||||
}
|
||||
// В endgame: убираем только этого пира из списка
|
||||
if peers, ok := endgamePeers[req.pieceIndex]; ok {
|
||||
filtered := peers[:0]
|
||||
for _, p := range peers {
|
||||
if p != "" { // убираем все (peerID недоступен в req)
|
||||
filtered = append(filtered, p)
|
||||
}
|
||||
}
|
||||
if len(filtered) == 0 {
|
||||
delete(endgamePeers, req.pieceIndex)
|
||||
} else {
|
||||
endgamePeers[req.pieceIndex] = filtered
|
||||
if states[req.pieceIndex] != pieceDone {
|
||||
states[req.pieceIndex] = piecePending
|
||||
}
|
||||
}
|
||||
req.responseCh <- false
|
||||
req.responseCh <- true
|
||||
|
||||
case mode := <-ps.setSequentialCh:
|
||||
sequentialMode = mode
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -448,6 +438,35 @@ func selectPendingPieceRarest(states []pieceState, availability []int, have []bo
|
|||
return -1
|
||||
}
|
||||
|
||||
func selectPendingPieceSequential(states []pieceState, have []bool, hasInfo bool) int {
|
||||
firstPending := firstPendingPiece(states)
|
||||
if firstPending < 0 {
|
||||
return -1
|
||||
}
|
||||
|
||||
// If peer has not sent bitfield/have yet, optimistically probe.
|
||||
if !hasInfo {
|
||||
return firstPending
|
||||
}
|
||||
|
||||
for pieceIndex, state := range states {
|
||||
if state != piecePending {
|
||||
continue
|
||||
}
|
||||
if pieceIndex >= len(have) || !have[pieceIndex] {
|
||||
continue
|
||||
}
|
||||
return pieceIndex // В последовательном режиме сразу берём первый доступный
|
||||
}
|
||||
|
||||
// Some peers send truncated availability info; allow fallback probing.
|
||||
if len(have) == 0 || len(have) < len(states) {
|
||||
return firstPending
|
||||
}
|
||||
|
||||
return -1
|
||||
}
|
||||
|
||||
func firstPendingPiece(states []pieceState) int {
|
||||
for pieceIndex, state := range states {
|
||||
if state == piecePending {
|
||||
|
|
|
|||
16
ui/ui.go
16
ui/ui.go
|
|
@ -122,6 +122,7 @@ type model struct {
|
|||
// Загрузка
|
||||
loadingMetadata bool
|
||||
spinner spinner.Model
|
||||
streamingMode bool
|
||||
|
||||
// Таблицы
|
||||
peersTable table.Model
|
||||
|
|
@ -470,6 +471,14 @@ func (m model) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
|||
m.activeTab = (m.activeTab + 1) % 4
|
||||
case key.Matches(msg, keys.PrevTab):
|
||||
m.activeTab = (m.activeTab - 1 + 4) % 4
|
||||
case msg.String() == "v" || msg.String() == "V":
|
||||
m.streamingMode = !m.streamingMode
|
||||
m.controller.SetSequentialMode(m.streamingMode)
|
||||
if m.streamingMode {
|
||||
m.controller.StartStreamServer()
|
||||
} else {
|
||||
m.controller.StopStreamServer()
|
||||
}
|
||||
case key.Matches(msg, keys.PauseResume):
|
||||
phase := m.status.Phase
|
||||
if phase == "stopped" || phase == "failed" || phase == "" {
|
||||
|
|
@ -1246,7 +1255,13 @@ func (m model) renderHeader() string {
|
|||
lipgloss.NewStyle().Foreground(lipgloss.Color(ColorText)).Render(name),
|
||||
)
|
||||
|
||||
var streamBadge string
|
||||
if m.streamingMode {
|
||||
streamBadge = lipgloss.NewStyle().Foreground(lipgloss.Color("#0b1017")).Background(lipgloss.Color(ColorPink)).Bold(true).Padding(0, 1).Render("▶ STREAM") + " "
|
||||
}
|
||||
|
||||
right := lipgloss.JoinHorizontal(lipgloss.Center,
|
||||
streamBadge,
|
||||
lipgloss.NewStyle().Foreground(lipgloss.Color(ColorSubdued)).Render(icon+" "),
|
||||
phaseColored,
|
||||
)
|
||||
|
|
@ -1644,6 +1659,7 @@ func (m model) renderHelpBar() string {
|
|||
k("J", "журнал"),
|
||||
k("L", "логи"),
|
||||
k("?", "инфо"),
|
||||
k("V", "стриминг"),
|
||||
k("Пробел", "пауза/старт"),
|
||||
k("O", "открыть"),
|
||||
k("C", "очистить"),
|
||||
|
|
|
|||
Loading…
Reference in a new issue