sheetsync: Replace old special-case PlaylistSync with SheetSync subclass

pull/401/head
Mike Lang 4 months ago committed by Mike Lang
parent ef11b69f4d
commit bb4cc8c668

@ -10,7 +10,6 @@ import gevent.event
import prometheus_client as prom import prometheus_client as prom
from monotonic import monotonic from monotonic import monotonic
from psycopg2 import sql from psycopg2 import sql
from psycopg2.extras import execute_values
from requests import HTTPError from requests import HTTPError
import common import common
@ -18,7 +17,7 @@ import common.dateutil
from common.database import DBManager, query, get_column_placeholder from common.database import DBManager, query, get_column_placeholder
from common.sheets import Sheets as SheetsClient from common.sheets import Sheets as SheetsClient
from .sheets import SheetsEventsMiddleware from .sheets import SheetsEventsMiddleware, SheetsPlaylistsMiddleware
from .streamlog import StreamLogClient, StreamLogMiddleware from .streamlog import StreamLogClient, StreamLogMiddleware
sheets_synced = prom.Counter( sheets_synced = prom.Counter(
@ -301,80 +300,20 @@ class EventsSync(SheetSync):
super().sync_row(sheet_row, db_row) super().sync_row(sheet_row, db_row)
class PlaylistSync:
# Time between syncs class PlaylistsSync(SheetSync):
RETRY_INTERVAL = 20
# Time to wait after getting an error
ERROR_RETRY_INTERVAL = 20
def __init__(self, stop, dbmanager, sheets, sheet_id, worksheet): # Slower poll rate than events to avoid using large amounts of quota
self.stop = stop retry_interval = 20
self.dbmanager = dbmanager error_retry_interval = 20
self.sheets = sheets
self.sheet_id = sheet_id
self.worksheet = worksheet
def run(self):
self.conn = self.dbmanager.get_conn()
while not self.stop.is_set(): table = "playlists"
try: input_columns = {
sync_start = monotonic() "playlist_id",
rows = self.sheets.get_rows(self.sheet_id, self.worksheet) "tags",
self.sync_playlists(rows) "name",
except Exception as e: "show_in_description",
# for HTTPErrors, http response body includes the more detailed error }
detail = ''
if isinstance(e, HTTPError):
detail = ": {}".format(e.response.content)
logging.exception("Failed to sync{}".format(detail))
sync_errors.labels("playlists").inc()
# To ensure a fresh slate and clear any DB-related errors, get a new conn on error.
# This is heavy-handed but simple and effective.
# If we can't re-connect, the program will crash from here,
# then restart and wait until it can connect again.
self.conn = self.dbmanager.get_conn()
wait(self.stop, sync_start, self.ERROR_RETRY_INTERVAL)
else:
logging.info("Successful sync of playlists")
sheets_synced.labels("playlists").inc()
sheet_sync_duration.labels("playlists").observe(monotonic() - sync_start)
wait(self.stop, sync_start, self.RETRY_INTERVAL)
def sync_playlists(self, rows):
"""Parse rows with a valid playlist id and at least one tag,
overwriting the entire playlists table"""
playlists = []
for row in rows:
if len(row) == 5:
tags, _, name, playlist_id, show_in_description = row
elif len(row) == 4:
tags, _, name, playlist_id = row
show_in_description = ""
else:
continue
tags = [tag.strip() for tag in tags.split(',') if tag.strip()]
if not tags:
continue
# special-case for the "all everything" list,
# we don't want "no tags" to mean "all videos" so we need a sentinel value.
if tags == ["<all>"]:
tags = []
playlist_id = playlist_id.strip()
if len(playlist_id) != 34 or not playlist_id.startswith('PL'):
continue
show_in_description = show_in_description == "[✓]"
playlists.append((playlist_id, tags, name, show_in_description))
# We want to wipe and replace all the current entries in the table.
# The easiest way to do this is a DELETE then an INSERT, all within a transaction.
# The "with" block will perform everything under it within a transaction, rolling back
# on error or committing on exit.
logging.info("Updating playlists table with {} playlists".format(len(playlists)))
with self.conn:
query(self.conn, "DELETE FROM playlists")
execute_values(self.conn.cursor(), "INSERT INTO playlists(playlist_id, tags, name, show_in_description) VALUES %s", playlists)
@argh.arg('dbconnect', help= @argh.arg('dbconnect', help=
@ -460,8 +399,14 @@ def main(dbconnect, sync_configs, metrics_port=8005, backdoor_port=0):
config.get("allocate_ids", False), config.get("allocate_ids", False),
) )
if "playlist_worksheet" in config: if "playlist_worksheet" in config:
middleware = SheetsPlaylistsMiddleware(
client,
config["sheet_id"],
[config["playlist_worksheet"]],
config.get("allocate_ids", False),
)
workers.append( workers.append(
PlaylistSync(stop, dbmanager, client, config["sheet_id"], config["playlist_worksheet"]) PlaylistsSync("playlists", middleware, stop, dbmanager, config.get("reverse_sync", False))
) )
elif config["type"] == "streamlog": elif config["type"] == "streamlog":
auth_token = open(config["creds"]).read().strip() auth_token = open(config["creds"]).read().strip()

Loading…
Cancel
Save