ytaudio: accurate active queue system
Replace the drained asyncio.Queue with an explicit, always-inspectable model: GuildPlayer.pending (waiting tracks) + current (now playing). - queue command shows Now playing plus numbered Up next list - play reports queue position when adding behind an active track - stop/leave use clear_pending(); skip relies on the TRACK_END event to advance exactly one track (no double-skip) - queue alias q, nowplaying alias np
This commit is contained in:
+53
-32
@@ -75,38 +75,58 @@ class QueuedTrack:
|
||||
|
||||
|
||||
class GuildPlayer:
|
||||
"""Per-guild queue + player-loop state. One per guild."""
|
||||
"""Per-guild queue + player-loop state. One per guild.
|
||||
|
||||
The queue model is explicit and always inspectable:
|
||||
* ``pending`` is the list of tracks waiting to play (index 0 is next).
|
||||
* ``current`` is the track playing right now (or None).
|
||||
Commands read these directly, so nowplaying/queue always reflect reality.
|
||||
"""
|
||||
|
||||
def __init__(self, cog: "YTAudio", guild_id: int):
|
||||
self.cog = cog
|
||||
self.guild_id = guild_id
|
||||
self.queue: "asyncio.Queue[Optional[QueuedTrack]]" = asyncio.Queue()
|
||||
self.pending: List[QueuedTrack] = []
|
||||
self.current: Optional[QueuedTrack] = None
|
||||
self.task: Optional[asyncio.Task] = None
|
||||
self.next_event = asyncio.Event()
|
||||
self.next_event = asyncio.Event() # fired when the current track ends
|
||||
self.wakeup = asyncio.Event() # fired when a track is added
|
||||
self.volume = 100
|
||||
|
||||
def start(self):
|
||||
if self.task is None or self.task.done():
|
||||
self.task = asyncio.create_task(self._player_loop())
|
||||
|
||||
def enqueue(self, track: "QueuedTrack"):
|
||||
self.pending.append(track)
|
||||
self.wakeup.set()
|
||||
|
||||
def clear_pending(self):
|
||||
cleared = self.pending
|
||||
self.pending = []
|
||||
for t in cleared:
|
||||
self.cog._cleanup_file(t)
|
||||
|
||||
async def _player_loop(self):
|
||||
while True:
|
||||
self.next_event.clear()
|
||||
if not self.pending:
|
||||
self.wakeup.clear()
|
||||
if not self.pending:
|
||||
try:
|
||||
track = await self.queue.get()
|
||||
await self.wakeup.wait()
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
if track is None: # shutdown sentinel
|
||||
return
|
||||
continue
|
||||
|
||||
track = self.pending.pop(0)
|
||||
self.current = track
|
||||
self.next_event.clear()
|
||||
try:
|
||||
await self.cog._play_track(self.guild_id, track, self.next_event)
|
||||
# Wait until the track-end handler fires next_event
|
||||
await self.next_event.wait()
|
||||
except asyncio.CancelledError:
|
||||
self.cog._cleanup_file(track)
|
||||
self.current = None
|
||||
return
|
||||
except Exception:
|
||||
log.exception("Error playing track in guild %s", self.guild_id)
|
||||
@@ -348,23 +368,26 @@ class YTAudio(commands.Cog):
|
||||
gp = self._get_gp(ctx.guild.id)
|
||||
gp.volume = await self.config.guild(ctx.guild).volume()
|
||||
track = QueuedTrack(query, ctx.author, ctx.channel)
|
||||
already_active = gp.current is not None or not gp.queue.empty()
|
||||
await gp.queue.put(track)
|
||||
already_active = gp.current is not None or bool(gp.pending)
|
||||
gp.enqueue(track)
|
||||
if already_active:
|
||||
await ctx.send("Added to queue: {}".format(query))
|
||||
position = len(gp.pending)
|
||||
await ctx.send("Added to queue (position {}): {}".format(position, query))
|
||||
|
||||
@commands.command()
|
||||
@commands.guild_only()
|
||||
async def skip(self, ctx: commands.Context):
|
||||
"""Skip the current track."""
|
||||
gp = self.players.get(ctx.guild.id)
|
||||
if gp is None or gp.current is None:
|
||||
return await ctx.send("Nothing is playing.")
|
||||
try:
|
||||
player = lavalink.get_player(ctx.guild.id)
|
||||
except Exception:
|
||||
return await ctx.send("Nothing is playing.")
|
||||
# Stopping fires TRACK_END, which the event handler uses to advance
|
||||
# the queue. Do not set next_event here or it would skip two tracks.
|
||||
await player.stop()
|
||||
gp = self.players.get(ctx.guild.id)
|
||||
if gp:
|
||||
gp.next_event.set()
|
||||
await ctx.send("Skipped.")
|
||||
|
||||
@commands.command()
|
||||
@@ -373,12 +396,7 @@ class YTAudio(commands.Cog):
|
||||
"""Stop playback and clear the queue."""
|
||||
gp = self.players.get(ctx.guild.id)
|
||||
if gp:
|
||||
while not gp.queue.empty():
|
||||
try:
|
||||
t = gp.queue.get_nowait()
|
||||
self._cleanup_file(t)
|
||||
except asyncio.QueueEmpty:
|
||||
break
|
||||
gp.clear_pending()
|
||||
gp.next_event.set()
|
||||
try:
|
||||
player = lavalink.get_player(ctx.guild.id)
|
||||
@@ -409,20 +427,27 @@ class YTAudio(commands.Cog):
|
||||
await player.pause(False)
|
||||
await ctx.send("Resumed.")
|
||||
|
||||
@commands.command(name="queue")
|
||||
@commands.command(name="queue", aliases=["q"])
|
||||
@commands.guild_only()
|
||||
async def queue_cmd(self, ctx: commands.Context):
|
||||
"""Show the queue."""
|
||||
"""Show the current track and the queue."""
|
||||
gp = self.players.get(ctx.guild.id)
|
||||
if not gp or (gp.current is None and gp.queue.empty()):
|
||||
if not gp or (gp.current is None and not gp.pending):
|
||||
return await ctx.send("The queue is empty.")
|
||||
lines = []
|
||||
if gp.current:
|
||||
lines.append("Now playing: {}".format(gp.current.title or gp.current.query))
|
||||
pending: List[QueuedTrack] = list(gp.queue._queue) # snapshot
|
||||
for i, t in enumerate(pending, 1):
|
||||
lines.append(f"{i}. {t.title or t.query}")
|
||||
await ctx.send("\n".join(lines[:20]))
|
||||
if gp.pending:
|
||||
lines.append("Up next:")
|
||||
for i, t in enumerate(gp.pending, 1):
|
||||
lines.append("{}. {}".format(i, t.title or t.query))
|
||||
else:
|
||||
lines.append("Queue is empty.")
|
||||
# Keep within Discord message limits.
|
||||
out = "\n".join(lines)
|
||||
if len(out) > 1900:
|
||||
out = out[:1900] + "\n..."
|
||||
await ctx.send(out)
|
||||
|
||||
@commands.command(name="nowplaying", aliases=["np"])
|
||||
@commands.guild_only()
|
||||
@@ -462,11 +487,7 @@ class YTAudio(commands.Cog):
|
||||
"""Leave the voice channel and clear state."""
|
||||
gp = self.players.get(ctx.guild.id)
|
||||
if gp:
|
||||
while not gp.queue.empty():
|
||||
try:
|
||||
self._cleanup_file(gp.queue.get_nowait())
|
||||
except asyncio.QueueEmpty:
|
||||
break
|
||||
gp.clear_pending()
|
||||
gp.next_event.set()
|
||||
try:
|
||||
player = lavalink.get_player(ctx.guild.id)
|
||||
|
||||
Reference in New Issue
Block a user