Merge pull request #356 from Dispatcharr:stats-direct-api

Switch to api calls instead of celery beat for stats
This commit is contained in:
SergeantPanda 2025-09-05 12:10:10 -05:00 committed by GitHub
commit b68b904838
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 303 additions and 91 deletions

View file

@ -12,6 +12,7 @@ from .constants import EventType
from .config_helper import ConfigHelper
from .redis_keys import RedisKeys
from .utils import get_logger
from core.utils import send_websocket_update
logger = get_logger()
@ -36,6 +37,45 @@ class ClientManager:
self._start_heartbeat_thread()
self._registered_clients = set() # Track already registered client IDs
def _trigger_stats_update(self):
"""Trigger a channel stats update via WebSocket"""
try:
# Import here to avoid potential import issues
from apps.proxy.ts_proxy.channel_status import ChannelStatus
import redis
# Get all channels from Redis
redis_client = redis.Redis.from_url('redis://localhost:6379', decode_responses=True)
all_channels = []
cursor = 0
while True:
cursor, keys = redis_client.scan(cursor, match="ts_proxy:channel:*:clients", count=100)
for key in keys:
# Extract channel ID from key
parts = key.split(':')
if len(parts) >= 4:
ch_id = parts[2]
channel_info = ChannelStatus.get_basic_channel_info(ch_id)
if channel_info:
all_channels.append(channel_info)
if cursor == 0:
break
# Send WebSocket update using existing infrastructure
send_websocket_update(
"updates",
"update",
{
"success": True,
"type": "channel_stats",
"stats": json.dumps({'channels': all_channels, 'count': len(all_channels)})
}
)
except Exception as e:
logger.debug(f"Failed to trigger stats update: {e}")
def _start_heartbeat_thread(self):
"""Start thread to regularly refresh client presence in Redis"""
def heartbeat_task():
@ -260,6 +300,9 @@ class ClientManager:
json.dumps(event_data)
)
# Trigger channel stats update via WebSocket
self._trigger_stats_update()
# Get total clients across all workers
total_clients = self.get_total_client_count()
logger.info(f"New client connected: {client_id} (local: {len(self.clients)}, total: {total_clients})")
@ -274,6 +317,8 @@ class ClientManager:
def remove_client(self, client_id):
"""Remove a client from this channel and Redis"""
client_ip = None
with self.lock:
if client_id in self.clients:
self.clients.remove(client_id)
@ -284,6 +329,14 @@ class ClientManager:
self.last_active_time = time.time()
if self.redis_client:
# Get client IP before removing the data
client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}"
client_data = self.redis_client.hgetall(client_key)
if client_data and b'ip_address' in client_data:
client_ip = client_data[b'ip_address'].decode('utf-8')
elif client_data and 'ip_address' in client_data:
client_ip = client_data['ip_address']
# Remove from channel's client set
self.redis_client.srem(self.client_set_key, client_id)
@ -313,6 +366,9 @@ class ClientManager:
})
self.redis_client.publish(RedisKeys.events_channel(self.channel_id), event_data)
# Trigger channel stats update via WebSocket
self._trigger_stats_update()
total_clients = self.get_total_client_count()
logger.info(f"Client disconnected: {client_id} (local: {len(self.clients)}, total: {total_clients})")

View file

@ -28,6 +28,7 @@ from apps.accounts.permissions import (
from .constants import ChannelState, EventType, StreamType, ChannelMetadataField
from .config_helper import ConfigHelper
from .services.channel_service import ChannelService
from core.utils import send_websocket_update
from .url_utils import (
generate_stream_url,
transform_url,
@ -633,6 +634,18 @@ def channel_status(request, channel_id=None):
if cursor == 0:
break
# Send WebSocket update with the stats
# Format it the same way the original Celery task did
send_websocket_update(
"updates",
"update",
{
"success": True,
"type": "channel_stats",
"stats": json.dumps({'channels': all_channels, 'count': len(all_channels)})
}
)
return JsonResponse({"channels": all_channels, "count": len(all_channels)})
except Exception as e:

View file

@ -198,10 +198,14 @@ CELERY_TASK_SERIALIZER = "json"
CELERY_BEAT_SCHEDULER = "django_celery_beat.schedulers.DatabaseScheduler"
CELERY_BEAT_SCHEDULE = {
# Explicitly disable the old fetch-channel-statuses task
# This ensures it gets disabled when DatabaseScheduler syncs
"fetch-channel-statuses": {
"task": "apps.proxy.tasks.fetch_channel_stats", # Direct task call
"schedule": 2.0, # Every 2 seconds
"task": "apps.proxy.tasks.fetch_channel_stats",
"schedule": 2.0, # Original schedule (doesn't matter since disabled)
"enabled": False, # Explicitly disabled
},
# Keep the file scanning task
"scan-files": {
"task": "core.tasks.scan_and_process_files", # Direct task call
"schedule": 20.0, # Every 20 seconds

View file

@ -194,7 +194,9 @@ export const WebsocketProvider = ({ children }) => {
loading: false,
autoClose: 4000,
});
try { await useChannelsStore.getState().fetchRecordings(); } catch {}
try {
await useChannelsStore.getState().fetchRecordings();
} catch {}
} else if (status === 'skipped') {
notifications.update({
id,
@ -204,7 +206,9 @@ export const WebsocketProvider = ({ children }) => {
loading: false,
autoClose: 3000,
});
try { await useChannelsStore.getState().fetchRecordings(); } catch {}
try {
await useChannelsStore.getState().fetchRecordings();
} catch {}
} else if (status === 'error') {
notifications.update({
id,
@ -214,7 +218,9 @@ export const WebsocketProvider = ({ children }) => {
loading: false,
autoClose: 6000,
});
try { await useChannelsStore.getState().fetchRecordings(); } catch {}
try {
await useChannelsStore.getState().fetchRecordings();
} catch {}
}
break;
}

View file

@ -1317,6 +1317,16 @@ export default class API {
}
}
static async fetchActiveChannelStats() {
try {
const response = await request(`${host}/proxy/ts/status`);
return response;
} catch (e) {
errorNotification('Failed to fetch active channel stats', e);
throw e;
}
}
static async getLogos(params = {}) {
try {
const queryParams = new URLSearchParams(params);

View file

@ -429,10 +429,18 @@ const SettingsPage = () => {
<Stack gap="sm">
<Switch
label="Enable Comskip (remove commercials after recording)"
{...form.getInputProps('dvr-comskip-enabled', { type: 'checkbox' })}
{...form.getInputProps('dvr-comskip-enabled', {
type: 'checkbox',
})}
key={form.key('dvr-comskip-enabled')}
id={settings['dvr-comskip-enabled']?.id || 'dvr-comskip-enabled'}
name={settings['dvr-comskip-enabled']?.key || 'dvr-comskip-enabled'}
id={
settings['dvr-comskip-enabled']?.id ||
'dvr-comskip-enabled'
}
name={
settings['dvr-comskip-enabled']?.key ||
'dvr-comskip-enabled'
}
/>
<TextInput
label="TV Path Template"
@ -440,8 +448,12 @@ const SettingsPage = () => {
placeholder="Recordings/TV_Shows/{show}/S{season:02d}E{episode:02d}.mkv"
{...form.getInputProps('dvr-tv-template')}
key={form.key('dvr-tv-template')}
id={settings['dvr-tv-template']?.id || 'dvr-tv-template'}
name={settings['dvr-tv-template']?.key || 'dvr-tv-template'}
id={
settings['dvr-tv-template']?.id || 'dvr-tv-template'
}
name={
settings['dvr-tv-template']?.key || 'dvr-tv-template'
}
/>
<TextInput
label="TV Fallback Template"
@ -449,8 +461,14 @@ const SettingsPage = () => {
placeholder="Recordings/TV_Shows/{show}/{start}.mkv"
{...form.getInputProps('dvr-tv-fallback-template')}
key={form.key('dvr-tv-fallback-template')}
id={settings['dvr-tv-fallback-template']?.id || 'dvr-tv-fallback-template'}
name={settings['dvr-tv-fallback-template']?.key || 'dvr-tv-fallback-template'}
id={
settings['dvr-tv-fallback-template']?.id ||
'dvr-tv-fallback-template'
}
name={
settings['dvr-tv-fallback-template']?.key ||
'dvr-tv-fallback-template'
}
/>
<TextInput
label="Movie Path Template"
@ -458,8 +476,14 @@ const SettingsPage = () => {
placeholder="Recordings/Movies/{title} ({year}).mkv"
{...form.getInputProps('dvr-movie-template')}
key={form.key('dvr-movie-template')}
id={settings['dvr-movie-template']?.id || 'dvr-movie-template'}
name={settings['dvr-movie-template']?.key || 'dvr-movie-template'}
id={
settings['dvr-movie-template']?.id ||
'dvr-movie-template'
}
name={
settings['dvr-movie-template']?.key ||
'dvr-movie-template'
}
/>
<TextInput
label="Movie Fallback Template"
@ -467,11 +491,24 @@ const SettingsPage = () => {
placeholder="Recordings/Movies/{start}.mkv"
{...form.getInputProps('dvr-movie-fallback-template')}
key={form.key('dvr-movie-fallback-template')}
id={settings['dvr-movie-fallback-template']?.id || 'dvr-movie-fallback-template'}
name={settings['dvr-movie-fallback-template']?.key || 'dvr-movie-fallback-template'}
id={
settings['dvr-movie-fallback-template']?.id ||
'dvr-movie-fallback-template'
}
name={
settings['dvr-movie-fallback-template']?.key ||
'dvr-movie-fallback-template'
}
/>
<Flex mih={50} gap="xs" justify="flex-end" align="flex-end">
<Button type="submit" variant="default">Save</Button>
<Flex
mih={50}
gap="xs"
justify="flex-end"
align="flex-end"
>
<Button type="submit" variant="default">
Save
</Button>
</Flex>
</Stack>
</form>

View file

@ -1,7 +1,8 @@
import React, { useMemo, useState, useEffect } from 'react';
import React, { useMemo, useState, useEffect, useCallback } from 'react';
import {
ActionIcon,
Box,
Button,
Card,
Center,
Container,
@ -12,9 +13,9 @@ import {
Text,
Title,
Tooltip,
useMantineTheme,
Select,
Badge,
NumberInput,
} from '@mantine/core';
import { TableHelper } from '../helpers';
import API from '../api';
@ -197,7 +198,7 @@ const ChannelCard = ({
...client,
}))
);
}, [clients]);
}, [clients, channel.channel_id]);
const renderHeaderCell = (header) => {
switch (header.id) {
@ -718,17 +719,23 @@ const ChannelCard = ({
};
const ChannelsPage = () => {
const theme = useMantineTheme();
const channels = useChannelsStore((s) => s.channels);
const channelsByUUID = useChannelsStore((s) => s.channelsByUUID);
const channelStats = useChannelsStore((s) => s.stats);
const logos = useLogosStore((s) => s.logos); // Add logos from the store
const setChannelStats = useChannelsStore((s) => s.setChannelStats);
const logos = useLogosStore((s) => s.logos);
const streamProfiles = useStreamProfilesStore((s) => s.profiles);
const fetchSettings = useSettingsStore((s) => s.fetchSettings);
const [activeChannels, setActiveChannels] = useState({});
const [clients, setClients] = useState([]);
const [isPollingActive, setIsPollingActive] = useState(false);
// Use localStorage for stats refresh interval (in seconds)
const [refreshIntervalSeconds, setRefreshIntervalSeconds] = useLocalStorage(
'stats-refresh-interval',
5
);
const refreshInterval = refreshIntervalSeconds * 1000; // Convert to milliseconds
const channelsColumns = useMemo(
() => [
@ -818,14 +825,57 @@ const ChannelsPage = () => {
await API.stopClient(channelId, clientId);
};
// The main clientsTable is no longer needed since each channel card has its own table
// Function to fetch channel stats from API
const fetchChannelStats = useCallback(async () => {
try {
const response = await API.fetchActiveChannelStats();
if (response) {
setChannelStats(response);
} else {
console.log('API response was empty or null');
}
} catch (error) {
console.error('Error fetching channel stats:', error);
console.error('Error details:', {
message: error.message,
status: error.status,
body: error.body,
});
}
}, [setChannelStats]);
// Fetch settings on component mount
// Set up polling for stats when on stats page
useEffect(() => {
fetchSettings();
}, [fetchSettings]);
const location = window.location;
const isOnStatsPage = location.pathname === '/stats';
if (isOnStatsPage && refreshInterval > 0) {
setIsPollingActive(true);
// Initial fetch
fetchChannelStats();
// Set up interval
const interval = setInterval(() => {
fetchChannelStats();
}, refreshInterval);
return () => {
clearInterval(interval);
setIsPollingActive(false);
};
} else {
setIsPollingActive(false);
}
}, [refreshInterval, fetchChannelStats]);
// Fetch initial stats on component mount (for immediate data when navigating to page)
useEffect(() => {
fetchChannelStats();
}, [fetchChannelStats]);
useEffect(() => {
console.log('Processing channel stats:', channelStats);
if (
!channelStats ||
!channelStats.channels ||
@ -834,81 +884,117 @@ const ChannelsPage = () => {
) {
console.log('No channel stats available:', channelStats);
// Clear active channels when there are no stats
if (Object.keys(activeChannels).length > 0) {
setActiveChannels({});
setClients([]);
}
setActiveChannels((prevActiveChannels) => {
if (Object.keys(prevActiveChannels).length > 0) {
setClients([]);
return {};
}
return prevActiveChannels;
});
return;
}
// Create a completely new object based only on current channel stats
const stats = {};
// Use functional update to access previous state without dependency
setActiveChannels((prevActiveChannels) => {
// Create a completely new object based only on current channel stats
const stats = {};
// Track which channels are currently active according to channelStats
const currentActiveChannelIds = new Set(
channelStats.channels.map((ch) => ch.channel_id).filter(Boolean)
);
channelStats.channels.forEach((ch) => {
// Make sure we have a valid channel_id
if (!ch.channel_id) {
console.warn('Found channel without channel_id:', ch);
return;
}
let bitrates = [];
if (activeChannels[ch.channel_id]) {
bitrates = [...(activeChannels[ch.channel_id].bitrates || [])];
const bitrate =
ch.total_bytes - activeChannels[ch.channel_id].total_bytes;
if (bitrate > 0) {
bitrates.push(bitrate);
channelStats.channels.forEach((ch) => {
// Make sure we have a valid channel_id
if (!ch.channel_id) {
console.warn('Found channel without channel_id:', ch);
return;
}
if (bitrates.length > 15) {
bitrates = bitrates.slice(1);
let bitrates = [];
if (prevActiveChannels[ch.channel_id]) {
bitrates = [...(prevActiveChannels[ch.channel_id].bitrates || [])];
const bitrate =
ch.total_bytes - prevActiveChannels[ch.channel_id].total_bytes;
if (bitrate > 0) {
bitrates.push(bitrate);
}
if (bitrates.length > 15) {
bitrates = bitrates.slice(1);
}
}
}
// Find corresponding channel data
const channelData =
channelsByUUID && ch.channel_id
? channels[channelsByUUID[ch.channel_id]]
: null;
// Find corresponding channel data
const channelData =
channelsByUUID && ch.channel_id
? channels[channelsByUUID[ch.channel_id]]
: null;
// Find stream profile
const streamProfile = streamProfiles.find(
(profile) => profile.id == parseInt(ch.stream_profile)
);
stats[ch.channel_id] = {
...ch,
...(channelData || {}), // Safely merge channel data if available
bitrates,
stream_profile: streamProfile || { name: 'Unknown' },
// Make sure stream_id is set from the active stream info
stream_id: ch.stream_id || null,
};
});
console.log('Processed active channels:', stats);
setActiveChannels(stats);
const clientStats = Object.values(stats).reduce((acc, ch) => {
if (ch.clients && Array.isArray(ch.clients)) {
return acc.concat(
ch.clients.map((client) => ({
...client,
channel: ch,
}))
// Find stream profile
const streamProfile = streamProfiles.find(
(profile) => profile.id == parseInt(ch.stream_profile)
);
}
return acc;
}, []);
setClients(clientStats);
stats[ch.channel_id] = {
...ch,
...(channelData || {}), // Safely merge channel data if available
bitrates,
stream_profile: streamProfile || { name: 'Unknown' },
// Make sure stream_id is set from the active stream info
stream_id: ch.stream_id || null,
};
});
console.log('Processed active channels:', stats);
// Update clients based on new stats
const clientStats = Object.values(stats).reduce((acc, ch) => {
if (ch.clients && Array.isArray(ch.clients)) {
return acc.concat(
ch.clients.map((client) => ({
...client,
channel: ch,
}))
);
}
return acc;
}, []);
setClients(clientStats);
return stats;
});
}, [channelStats, channels, channelsByUUID, streamProfiles]);
return (
<Box style={{ overflowX: 'auto' }}>
<Box style={{ padding: '10px', borderBottom: '1px solid #444' }}>
<Group justify="space-between" align="center">
<Title order={3}>Active Channels</Title>
<Group align="center">
<NumberInput
label="Refresh Interval (seconds)"
value={refreshIntervalSeconds}
onChange={(value) => setRefreshIntervalSeconds(value || 0)}
min={0}
max={300}
step={1}
size="xs"
style={{ width: 120 }}
description={
refreshIntervalSeconds === 0 ? 'Disabled' : 'Auto-refresh'
}
/>
{isPollingActive && refreshInterval > 0 && (
<Text size="sm" c="dimmed">
Refreshing every {refreshIntervalSeconds}s
</Text>
)}
<Button
size="xs"
variant="subtle"
onClick={fetchChannelStats}
loading={false}
>
Refresh Now
</Button>
</Group>
</Group>
</Box>
<div
style={{
display: 'grid',