mirror of
https://github.com/JohnDoee/deluge-streaming/
synced 2026-07-01 07:31:17 -07:00
Merge branch 'release/0.7.0'
This commit is contained in:
10
README.md
10
README.md
@@ -1,7 +1,7 @@
|
||||
# Streaming Plugin
|
||||
https://github.com/JohnDoee/deluge-streaming
|
||||
|
||||
(c)2015 by Anders Jensen <johndoee@tidalstream.org>
|
||||
(c)2016 by Anders Jensen <johndoee@tidalstream.org>
|
||||
|
||||
## Description
|
||||
|
||||
@@ -35,12 +35,16 @@ make Deluge an abstraction layer for the [TidalStream](http://www.tidalstream.or
|
||||
|
||||
The _allow remote_ option is to allow remote add and stream of torrents.
|
||||
|
||||
## ToDo
|
||||
# Known bugs
|
||||
|
||||
* Add feedback when preparing stream.
|
||||
* Sometimes the plugin tries to read empty data when there is too much requesting going on.
|
||||
|
||||
# Version Info
|
||||
|
||||
## Version 0.7.0
|
||||
* Shrinked code by redoing queue algorithm. This should prevent more stalled downloads and allow it to act bittorrenty if necessary.
|
||||
* Added support for waiting for end pieces to satisfy some video players (KODI)
|
||||
|
||||
## Version 0.6.1
|
||||
* Should not have been in changelog: Fixed "resume on complete" broken-ness (i hope)
|
||||
|
||||
|
||||
3
setup.py
3
setup.py
@@ -42,7 +42,7 @@ from setuptools import setup
|
||||
__plugin_name__ = "Streaming"
|
||||
__author__ = "John Doee"
|
||||
__author_email__ = "johndoee@tidalstream.org"
|
||||
__version__ = "0.6.1"
|
||||
__version__ = "0.7.0"
|
||||
__url__ = "https://github.com/JohnDoee/deluge-streaming"
|
||||
__license__ = "GPLv3"
|
||||
__description__ = "Enables streaming of files while downloading them."
|
||||
@@ -59,7 +59,6 @@ downloads ahead, this enables seeking in video files.
|
||||
* Select _files_ tab
|
||||
* Right-click a file.
|
||||
* Click _Stream this file_
|
||||
* **WAIT**, it will try to buffer the first pieces of the file before generating a link (no feedback yet).
|
||||
* Select the link, open it in a media player, e.g. VLC or MPC
|
||||
|
||||
If you want to stream from a non-local computer, e.g. your seedbox, you will need to change the IP in option to the external server ip."""
|
||||
|
||||
@@ -72,6 +72,9 @@ DEFAULT_PREFS = {
|
||||
'remote_password': 'password',
|
||||
}
|
||||
|
||||
# TODO: set priority for all torrents we're streaming using set_priority
|
||||
|
||||
PRIORITY_INCREASE = 5
|
||||
|
||||
def sleep(seconds):
|
||||
d = defer.Deferred()
|
||||
@@ -182,6 +185,7 @@ class TorrentFile(object): # can be read from, knows about itself
|
||||
self.index = index
|
||||
|
||||
self.file_requested = False
|
||||
self.file_requested_once = False
|
||||
self.do_shutdown = False
|
||||
self.first_piece_end = self.piece_size * (self.first_piece + 1) - offset
|
||||
self.waiting_pieces = {}
|
||||
@@ -198,17 +202,15 @@ class TorrentFile(object): # can be read from, knows about itself
|
||||
if not self.registered_alert:
|
||||
self.alerts.register_handler("read_piece_alert", self.on_alert_got_piece_data)
|
||||
self.registered_alert = True
|
||||
|
||||
tfr = TorrentFileReader(self)
|
||||
self.current_readers.append(tfr)
|
||||
self.file_requested = False
|
||||
|
||||
return tfr
|
||||
|
||||
def close(self, tfr):
|
||||
self.current_readers.remove(tfr)
|
||||
self.torrent.unprioritize_pieces(tfr)
|
||||
|
||||
if not self.current_readers:
|
||||
self.torrent.unprioritize_pieces(self)
|
||||
|
||||
def is_complete(self):
|
||||
torrent_status = self.torrent.torrent.get_status(['file_progress', 'state'])
|
||||
@@ -227,10 +229,25 @@ class TorrentFile(object): # can be read from, knows about itself
|
||||
if alert.piece not in self.waiting_pieces:
|
||||
logger.debug('Got data for piece %i, but no data needed for this piece?' % alert.piece)
|
||||
return
|
||||
# TODO: check piece size is not zero
|
||||
|
||||
if alert.buffer is None:
|
||||
return
|
||||
|
||||
self.waiting_pieces[alert.piece].callback(alert.buffer)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def wait_for_end_pieces(self):
|
||||
handle = self.torrent.torrent.handle
|
||||
for piece in [self.first_piece, self.last_piece]:
|
||||
handle.set_piece_deadline(piece, 0)
|
||||
handle.piece_priority(piece, 7)
|
||||
|
||||
while not handle.have_piece(self.first_piece) and not handle.have_piece(self.last_piece):
|
||||
if self.do_shutdown:
|
||||
raise Exception('Shutting down')
|
||||
logger.debug('Did not have piece %i, waiting' % piece)
|
||||
yield sleep(1)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def get_piece_data(self, piece):
|
||||
logger.debug('Trying to get piece data for piece %s' % piece)
|
||||
@@ -242,12 +259,10 @@ class TorrentFile(object): # can be read from, knows about itself
|
||||
self.waiting_pieces[piece] = defer.Deferred()
|
||||
|
||||
logger.debug('Waiting for %s' % piece)
|
||||
self.torrent.schedule_piece(self, self, piece, 0)
|
||||
while not self.torrent.torrent.handle.have_piece(piece):
|
||||
if self.do_shutdown:
|
||||
raise Exception()
|
||||
raise Exception('Shutting down')
|
||||
logger.debug('Did not have piece %i, waiting' % piece)
|
||||
self.torrent.unrelease()
|
||||
yield sleep(1)
|
||||
|
||||
self.torrent.torrent.handle.read_piece(piece)
|
||||
@@ -281,15 +296,8 @@ class Torrent(object):
|
||||
|
||||
self.last_piece = self.torrent_files[-1].last_piece
|
||||
self.torrent.handle.set_sequential_download(True)
|
||||
self.torrent.handle.set_priority(1)
|
||||
reactor.callLater(0, self.update_piece_priority)
|
||||
reactor.callLater(0, self.blackhole_all_pieces, 0, self.last_piece)
|
||||
|
||||
def is_unstreamed(self): # check if anyone streams this file
|
||||
for tf in self.torrent_files:
|
||||
if tf.current_readers:
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
def populate_files(self):
|
||||
self.torrent_files = []
|
||||
@@ -337,22 +345,6 @@ class Torrent(object):
|
||||
|
||||
return f
|
||||
|
||||
def unprioritize_pieces(self, tfr):
|
||||
logger.debug('Unprioritizing pieces for %s' % tfr)
|
||||
|
||||
if self.torrent_handler.config['reset_complete'] and self.is_unstreamed():
|
||||
logger.debug('Not unprioritizing pieces because there are no streams and it would stop the file (causes problems)')
|
||||
return
|
||||
|
||||
currently_downloading = self.get_currently_downloading()
|
||||
|
||||
for piece, increased_by in self.priority_increased.items():
|
||||
if tfr in increased_by:
|
||||
increased_by.remove(tfr)
|
||||
if not increased_by and piece not in currently_downloading and not self.torrent.status.pieces[piece]:
|
||||
logger.debug('Unprioritizing piece %s' % piece)
|
||||
self.torrent.handle.piece_priority(piece, 0)
|
||||
|
||||
def get_currently_downloading(self):
|
||||
currently_downloading = set()
|
||||
for peer in self.torrent.handle.get_peer_info():
|
||||
@@ -361,27 +353,10 @@ class Torrent(object):
|
||||
|
||||
return currently_downloading
|
||||
|
||||
def blackhole_all_pieces(self, first_piece, last_piece):
|
||||
currently_downloading = self.get_currently_downloading()
|
||||
|
||||
logger.debug('Blacklisting pieces from %i to %i skipping %r' % (first_piece, last_piece, currently_downloading))
|
||||
for piece in range(first_piece, last_piece+1):
|
||||
if piece not in currently_downloading and not self.torrent.status.pieces[piece]:
|
||||
if piece in self.priority_increased:
|
||||
continue
|
||||
#logger.debug('Setting piece priority %s to blacklist' % piece)
|
||||
self.torrent.handle.piece_priority(piece, 0)
|
||||
|
||||
def unrelease(self):
|
||||
if self.torrent_released:
|
||||
logger.debug('Unreleasing %s' % self.infohash)
|
||||
self.torrent_released = False
|
||||
self.torrent.set_file_priorities(self.file_priorities)
|
||||
self.blackhole_all_pieces(0, self.last_piece)
|
||||
|
||||
def get_torrent_file(self, file_or_index, includes_name):
|
||||
f = self.get_file(file_or_index, includes_name)
|
||||
f.file_requested = True
|
||||
f.file_requested_once = True
|
||||
|
||||
self.torrent.resume()
|
||||
|
||||
@@ -394,114 +369,124 @@ class Torrent(object):
|
||||
should_update_priorities = True
|
||||
|
||||
if should_update_priorities and not f.is_complete(): # Need to do this stuff on seek too
|
||||
self.unrelease()
|
||||
self.torrent.set_file_priorities(self.file_priorities)
|
||||
|
||||
return f
|
||||
|
||||
def shutdown(self):
|
||||
logger.info('Shutting down torrent %s' % self.infohash)
|
||||
logger.info('Shutting down torrent %s' % (self.infohash, ))
|
||||
|
||||
self.torrent.handle.set_priority(0)
|
||||
|
||||
for piece, status in enumerate(self.torrent.status.pieces[0:self.last_piece+1]):
|
||||
if status:
|
||||
continue
|
||||
|
||||
priority = self.torrent.handle.piece_priority(piece)
|
||||
if priority == 0:
|
||||
self.torrent.handle.piece_priority(piece, 1)
|
||||
|
||||
if self.torrent_handler.config['reset_complete']:
|
||||
logger.debug('Resetting file priorities')
|
||||
file_priorities = [(1 if fp == 0 else fp) for fp in self.file_priorities]
|
||||
self.torrent.set_file_priorities(file_priorities)
|
||||
|
||||
self.do_shutdown = True
|
||||
self.torrent_handler.remove_torrent(self.infohash)
|
||||
|
||||
for tf in self.torrent_files:
|
||||
tf.shutdown()
|
||||
|
||||
def schedule_piece(self, torrent_file, schedule_target, piece, distance):
|
||||
if schedule_target not in self.priority_increased[piece]:
|
||||
if not self.priority_increased[piece]:
|
||||
logger.debug('Scheduled piece %s at distance %s' % (piece, distance))
|
||||
|
||||
self.torrent.handle.piece_priority(piece, (7 if distance <= 4 else 6))
|
||||
self.torrent.handle.set_piece_deadline(piece, 700*(distance+1))
|
||||
self.priority_increased[piece].add(schedule_target)
|
||||
|
||||
def do_pieces_schedule(self, torrent_file, schedule_target, currently_downloading, from_piece):
|
||||
logger.debug('Looking for stuff to do with pieces for file %s from piece %s' % (torrent_file, from_piece))
|
||||
|
||||
priority_increased = 0
|
||||
chain_size = 0
|
||||
download_chain_size = 0
|
||||
end_of_chain = False
|
||||
|
||||
current_buffer_offset = 5
|
||||
if self.torrent.status.pieces[torrent_file.last_piece]:
|
||||
if torrent_file.first_piece != torrent_file.last_piece and self.torrent.status.pieces[torrent_file.last_piece-1] \
|
||||
or torrent_file.first_piece == torrent_file.last_piece:
|
||||
current_buffer_offset = 20
|
||||
|
||||
for piece, status in enumerate(self.torrent.status.pieces[from_piece:torrent_file.last_piece+1], from_piece):
|
||||
if not end_of_chain:
|
||||
if status:
|
||||
chain_size += 1
|
||||
elif piece in currently_downloading:
|
||||
download_chain_size += 1
|
||||
|
||||
if not status and piece not in currently_downloading:
|
||||
if not end_of_chain:
|
||||
status_increase = max(11, chain_size-current_buffer_offset)
|
||||
end_of_chain = True
|
||||
|
||||
priority_increased += 1
|
||||
if priority_increased >= status_increase:
|
||||
logger.debug('Done increasing priority for %i pieces' % status_increase)
|
||||
break
|
||||
|
||||
self.schedule_piece(torrent_file, schedule_target, piece, piece-from_piece)
|
||||
else:
|
||||
logger.info('We are done with the rest of this chain, we might be able to increase others')
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
def update_piece_priority(self): # if all do_pieces_schedule returns true, allow all pices of file to be downloaded or whole torernt
|
||||
def update_piece_priority(self): # if file streamed has reached end, unblacklist all prior pieces
|
||||
if self.do_shutdown:
|
||||
return
|
||||
|
||||
logger.debug('Updating piece priority for %s' % self.infohash)
|
||||
logger.debug('Updating piece priority for %s' % (self.infohash, ))
|
||||
currently_downloading = self.get_currently_downloading()
|
||||
|
||||
currently_downloading = set()
|
||||
for peer in self.torrent.handle.get_peer_info():
|
||||
if peer.downloading_piece_index != -1:
|
||||
currently_downloading.add(peer.downloading_piece_index)
|
||||
|
||||
all_heads_done = True
|
||||
for f in self.torrent_files:
|
||||
if not f.file_requested and not f.current_readers:
|
||||
if not f.file_requested and not f.current_readers: # nobody wants the file and nobody is watching
|
||||
continue
|
||||
|
||||
logger.debug('Rescheduling file %s' % f.path)
|
||||
logger.debug('Rescheduling file %s' % (f.path, ))
|
||||
|
||||
if f.file_requested:
|
||||
all_heads_done &= self.do_pieces_schedule(f, f, currently_downloading, f.first_piece)
|
||||
self.schedule_piece(f, f, f.last_piece, 0)
|
||||
if f.first_piece != f.last_piece:
|
||||
self.schedule_piece(f, f, f.last_piece-1, 1)
|
||||
heads = set()
|
||||
if f.file_requested: # we expect a piece head to be at start
|
||||
heads.add(f.first_piece)
|
||||
|
||||
waiting_for_pieces = set()
|
||||
|
||||
for tfr in f.current_readers:
|
||||
if tfr.waiting_for_piece is not None:
|
||||
logger.debug('Scheduling based on waiting for piece %s' % tfr.waiting_for_piece)
|
||||
all_heads_done &= self.do_pieces_schedule(f, tfr, currently_downloading, tfr.waiting_for_piece)
|
||||
elif tfr.current_piece is not None:
|
||||
logger.debug('Scheduling based on current piece %s' % tfr.current_piece)
|
||||
all_heads_done &= self.do_pieces_schedule(f, tfr, currently_downloading, tfr.current_piece)
|
||||
|
||||
if all(self.torrent.status.pieces):
|
||||
logger.debug('All pieces complete, no need to loop')
|
||||
return
|
||||
|
||||
if all_heads_done and not self.torrent_released:
|
||||
logger.debug('We are already done with all heads, figuring out what to do next')
|
||||
waiting_for_pieces.add(tfr.waiting_for_piece)
|
||||
|
||||
piece = max(tfr.waiting_for_piece, tfr.current_piece)
|
||||
if piece is not None:
|
||||
heads.add(piece)
|
||||
|
||||
if self.torrent_handler.config['reset_complete']:
|
||||
self.torrent_released = True
|
||||
logger.debug('Resetting all disabled files')
|
||||
file_priorities = [(1 if fp == 0 else fp) for fp in self.file_priorities]
|
||||
self.torrent.set_file_priorities(file_priorities)
|
||||
if not heads:
|
||||
continue
|
||||
|
||||
first_head = min(heads)
|
||||
|
||||
for head_piece in heads:
|
||||
priority_increased = 0
|
||||
for piece, status in enumerate(self.torrent.status.pieces[head_piece:f.last_piece+1], head_piece):
|
||||
#logger.debug('Checking status for %s/%s/%s/%s' % (head_piece, piece, status, self.torrent.handle.piece_priority(piece)))
|
||||
if status or piece in currently_downloading:
|
||||
continue
|
||||
|
||||
priority = self.torrent.handle.piece_priority(piece)
|
||||
if priority_increased < PRIORITY_INCREASE:
|
||||
priority_increased += 1
|
||||
|
||||
if piece in waiting_for_pieces:
|
||||
if priority < 7:
|
||||
logger.debug('setting priority for %s to 7 with deadline 0' % (piece, ))
|
||||
|
||||
self.torrent.handle.set_piece_deadline(piece, 0)
|
||||
self.torrent.handle.piece_priority(piece, 7)
|
||||
elif priority < 6:
|
||||
deadline = 3000 * priority_increased
|
||||
logger.debug('setting priority for %s to 6 with deadline %s' % (piece, deadline, ))
|
||||
self.torrent.handle.piece_priority(piece, 6)
|
||||
self.torrent.handle.set_piece_deadline(piece, deadline)
|
||||
|
||||
elif priority == 0:
|
||||
self.torrent.handle.piece_priority(piece, 1)
|
||||
|
||||
if head_piece == first_head:
|
||||
if priority_increased < PRIORITY_INCREASE:
|
||||
logger.debug('Everything we need has been scheduled, looking for pieces across file to unblacklist')
|
||||
for piece, status in enumerate(self.torrent.status.pieces[f.first_piece:f.last_piece+1], f.first_piece):
|
||||
if status:
|
||||
continue
|
||||
|
||||
priority = self.torrent.handle.piece_priority(piece)
|
||||
if priority == 0:
|
||||
self.torrent.handle.piece_priority(piece, 1)
|
||||
else:
|
||||
logger.debug('Looking for pieces before smallest head %s to blacklist' % (first_head, ))
|
||||
for piece, status in enumerate(self.torrent.status.pieces[f.first_piece:first_head], f.first_piece):
|
||||
if status or piece in currently_downloading:
|
||||
continue
|
||||
|
||||
if self.torrent.handle.piece_priority(piece) != 0:
|
||||
logger.debug('Blacklisting %i' % (piece, ))
|
||||
self.torrent.handle.piece_priority(piece, 0)
|
||||
|
||||
if not all_heads_done and self.torrent_released:
|
||||
logger.debug('Seems like the torrent was released too early')
|
||||
self.unrelease()
|
||||
found_requested = False
|
||||
for f in self.torrent_files:
|
||||
if f.file_requested_once:
|
||||
found_requested = True
|
||||
if not f.is_complete() or f.current_readers:
|
||||
break
|
||||
else:
|
||||
if found_requested:
|
||||
logger.debug('Nobody is currently using %s, shutting down torrent-handler' % (self.infohash, ))
|
||||
self.shutdown()
|
||||
|
||||
reactor.callLater(0.3, self.update_piece_priority)
|
||||
reactor.callLater(1, self.update_piece_priority)
|
||||
|
||||
class TorrentHandler(object):
|
||||
def __init__(self, config):
|
||||
@@ -529,6 +514,9 @@ class TorrentHandler(object):
|
||||
return
|
||||
|
||||
self.torrents[torrent_id].shutdown()
|
||||
self.remove_torrent(torrent_id)
|
||||
|
||||
def remove_torrent(self, torrent_id):
|
||||
del self.torrents[torrent_id]
|
||||
|
||||
def shutdown(self):
|
||||
@@ -593,7 +581,7 @@ class Core(CorePluginBase):
|
||||
|
||||
@export
|
||||
@defer.inlineCallbacks
|
||||
def stream_torrent(self, infohash=None, url=None, filedump=None, filepath_or_index=None, includes_name=False):
|
||||
def stream_torrent(self, infohash=None, url=None, filedump=None, filepath_or_index=None, includes_name=False, wait_for_end_pieces=False):
|
||||
tor = component.get("TorrentManager").torrents.get(infohash, None)
|
||||
|
||||
if tor is None:
|
||||
@@ -619,6 +607,10 @@ class Core(CorePluginBase):
|
||||
except UnknownTorrentException:
|
||||
defer.returnValue({'status': 'error', 'message': 'unable to find torrent, probably failed to add it'})
|
||||
|
||||
if wait_for_end_pieces:
|
||||
logger.debug('Waiting for end pieces')
|
||||
yield tf.wait_for_end_pieces()
|
||||
|
||||
defer.returnValue({
|
||||
'status': 'success',
|
||||
'use_stream_urls': self.config['use_stream_urls'],
|
||||
|
||||
Reference in New Issue
Block a user