diff --git a/internal/app/controller.go b/internal/app/controller.go index 7311b0f..71392ae 100644 --- a/internal/app/controller.go +++ b/internal/app/controller.go @@ -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 { @@ -22,9 +25,10 @@ func NewController() *Controller { ctx, cancel := context.WithCancel(context.Background()) return &Controller{ - engine: engine, - ctx: ctx, - cancel: cancel, + 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() +} diff --git a/internal/stream/server.go b/internal/stream/server.go new file mode 100644 index 0000000..421dc73 --- /dev/null +++ b/internal/stream/server.go @@ -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 +} diff --git a/internal/torrent/engine.go b/internal/torrent/engine.go index 69638c8..5dd032e 100644 --- a/internal/torrent/engine.go +++ b/internal/torrent/engine.go @@ -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 +} diff --git a/internal/torrent/piecescheduler.go b/internal/torrent/piecescheduler.go index 4212148..0a00aaa 100644 --- a/internal/torrent/piecescheduler.go +++ b/internal/torrent/piecescheduler.go @@ -42,6 +42,8 @@ type pieceScheduler struct { stopCh chan struct{} stopOnce sync.Once + setSequentialCh chan bool + // cancelCh рассылает сигнал отмены конкретному воркеру в endgame-режиме. // Ключ — peerID, значение — канал с индексом куска для отмены. cancelMu sync.RWMutex @@ -65,13 +67,14 @@ func newPieceScheduler(pieceCount int) *pieceScheduler { func newPieceSchedulerWithResume(pieceCount int, completedIndices []int) *pieceScheduler { ps := &pieceScheduler{ - assignCh: make(chan assignPieceRequest, 128), - reportCh: make(chan reportPieceRequest, 128), - releasePeerCh: make(chan string, 128), - progressCh: make(chan int, 128), - doneCh: make(chan struct{}), - stopCh: make(chan struct{}), - cancelSubs: make(map[string]chan int), + assignCh: make(chan assignPieceRequest, 128), + reportCh: make(chan reportPieceRequest, 128), + releasePeerCh: make(chan string, 128), + progressCh: make(chan int, 1), + doneCh: make(chan struct{}), + stopCh: make(chan struct{}), + setSequentialCh: make(chan bool), + cancelSubs: make(map[string]chan int), } go ps.run(pieceCount, completedIndices) @@ -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 { - pieceIndex = selectPendingPieceRarest(states, availability, req.have, req.hasInfo) + 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 + } else { + if states[req.pieceIndex] != pieceDone { + states[req.pieceIndex] = piecePending + } } + req.responseCh <- true - // Неуспех: возвращаем кусок в 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 - } - } - req.responseCh <- false + 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 { diff --git a/ui/ui.go b/ui/ui.go index 5745276..5ed05c7 100644 --- a/ui/ui.go +++ b/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", "очистить"),