Files
MoviePilot/app/db/models/transferhistory.py
T

560 lines
20 KiB
Python

import re
import time
from pathlib import Path
from typing import List, Optional
from sqlalchemy import Boolean, Column, Index, Integer, JSON, String, func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import Session
from app.db import db_query, db_update, get_id_column, Base, async_db_query
from app.schemas.types import MediaType
def _text_like(column, pattern: str, wildcard: bool = False):
"""构造跨数据库大小写不敏感的文本匹配条件。"""
if wildcard:
return column.ilike(pattern, escape='\\')
return column.ilike(pattern)
class TransferHistory(Base):
"""
整理记录
"""
id = get_id_column()
# 源路径
src = Column(String, index=True)
# 源存储
src_storage = Column(String)
# 源文件项
src_fileitem = Column(JSON, default=dict)
# 目标路径
dest = Column(String)
# 目标存储
dest_storage = Column(String)
# 目标文件项
dest_fileitem = Column(JSON, default=dict)
# 转移模式 move/copy/link...
mode = Column(String)
# 类型 电影/电视剧
type = Column(String)
# 二级分类
category = Column(String)
# 标题
title = Column(String, index=True)
# 年份
year = Column(String)
tmdbid = Column(Integer, index=True)
imdbid = Column(String)
tvdbid = Column(Integer)
doubanid = Column(String)
bangumiid = Column(Integer, index=True)
anilistid = Column(Integer, index=True)
# 统一媒体数据源与原生ID
media_source = Column(String, index=True)
media_id = Column(String, index=True)
# Sxx
seasons = Column(String)
# Exx
episodes = Column(String)
# 海报
image = Column(String)
# 下载器
downloader = Column(String)
# 下载器hash
download_hash = Column(String, index=True)
# 转移成功状态
status = Column(Boolean(), default=True)
# 转移失败信息
errmsg = Column(String)
# 时间
date = Column(String)
# 文件清单,以JSON存储
files = Column(JSON, default=list)
# 剧集组
episode_group = Column(String)
__table_args__ = (
Index('ix_transferhistory_status_date', 'status', 'date'),
Index('ix_transferhistory_date_id', 'date', 'id'),
Index('ix_transferhistory_media_identity', 'media_source', 'media_id'),
)
@classmethod
@db_query
def list_by_title(cls, db: Session, title: str, page: Optional[int] = 1, count: Optional[int] = 30,
status: bool = None, wildcard: bool = False):
if wildcard:
text_filter = or_(
_text_like(cls.title, title, wildcard=True),
_text_like(cls.src, title, wildcard=True),
_text_like(cls.dest, title, wildcard=True),
)
else:
text_filter = or_(
_text_like(cls.title, f'%{title}%'),
_text_like(cls.src, f'%{title}%'),
_text_like(cls.dest, f'%{title}%'),
)
query = db.query(cls).filter(text_filter)
if status is not None:
query = query.filter(cls.status == status)
query = query.order_by(cls.date.desc())
# 当count为负数时,不限制页数查询所有
if count >= 0:
query = query.offset((page - 1) * count).limit(count)
return query.all()
@classmethod
@async_db_query
async def async_list_by_title(cls, db: AsyncSession, title: str, page: Optional[int] = 1, count: Optional[int] = 30,
status: bool = None, wildcard: bool = False):
if wildcard:
text_filter = or_(
_text_like(cls.title, title, wildcard=True),
_text_like(cls.src, title, wildcard=True),
_text_like(cls.dest, title, wildcard=True),
)
else:
text_filter = or_(
_text_like(cls.title, f'%{title}%'),
_text_like(cls.src, f'%{title}%'),
_text_like(cls.dest, f'%{title}%'),
)
query = select(cls).filter(text_filter)
if status is not None:
query = query.filter(cls.status == status)
query = query.order_by(cls.date.desc())
# 当count为负数时,不限制页数查询所有
if count >= 0:
query = query.offset((page - 1) * count).limit(count)
result = await db.execute(query)
return result.scalars().all()
@classmethod
@db_query
def list_by_page(cls, db: Session, page: Optional[int] = 1, count: Optional[int] = 30, status: bool = None):
if status is not None:
query = db.query(cls).filter(
cls.status == status
).order_by(
cls.date.desc()
)
else:
query = db.query(cls).order_by(
cls.date.desc()
)
# 当count为负数时,不限制页数查询所有
if count >= 0:
query = query.offset((page - 1) * count).limit(count)
return query.all()
@classmethod
@async_db_query
async def async_list_by_page(cls, db: AsyncSession, page: Optional[int] = 1, count: Optional[int] = 30,
status: bool = None):
if status is not None:
query = select(cls).filter(
cls.status == status
).order_by(
cls.date.desc()
)
else:
query = select(cls).order_by(
cls.date.desc()
)
# 当count为负数时,不限制页数查询所有
if count >= 0:
query = query.offset((page - 1) * count).limit(count)
result = await db.execute(query)
return result.scalars().all()
@classmethod
@db_query
def get_by_hash(cls, db: Session, download_hash: str):
return db.query(cls).filter(cls.download_hash == download_hash).first()
@classmethod
@db_query
def get_by_src(
cls, db: Session, src: str, storage: Optional[str] = None
) -> Optional["TransferHistory"]:
"""
按源路径和存储查询单条整理记录。
:param db: 数据库会话
:param src: 源路径
:param storage: 源存储类型
:return: 命中的整理记录,未命中时返回 None
"""
if storage:
return db.query(cls).filter(cls.src == src,
cls.src_storage == storage).first()
else:
return db.query(cls).filter(cls.src == src).first()
@classmethod
@db_query
def get_by_dest(
cls, db: Session, dest: str, storage: Optional[str] = None
) -> Optional["TransferHistory"]:
"""
按目标路径和存储查询单条整理记录。
:param db: 数据库会话
:param dest: 目标路径
:param storage: 目标存储类型
:return: 命中的整理记录,未命中时返回 None
"""
query = db.query(cls).filter(cls.dest == dest)
if storage:
query = query.filter(cls.dest_storage == storage)
return query.first()
@classmethod
@db_query
def list_success_by_src(
cls,
db: Session,
src: str,
storage: Optional[str] = None,
recursive: bool = False,
) -> List["TransferHistory"]:
"""
按源路径查询成功整理记录,目录模式仅匹配其直接或间接子项。
:param db: 数据库会话
:param src: 源路径
:param storage: 源存储类型
:param recursive: 是否递归匹配目录子项
:return: 命中的成功整理记录
"""
normalized_src = (
Path(str(src).replace("\\", "/")).as_posix().rstrip("/") or "/"
)
query = db.query(cls).filter(cls.status.is_(True))
if recursive:
escaped_src = (
normalized_src.replace("\\", "\\\\")
.replace("%", "\\%")
.replace("_", "\\_")
)
query = query.filter(
or_(
cls.src == normalized_src,
cls.src.like(f"{escaped_src.rstrip('/')}/%", escape="\\"),
)
)
else:
query = query.filter(cls.src == normalized_src)
if storage:
query = query.filter(cls.src_storage == storage)
return query.all()
@classmethod
@db_query
def list_success_move_by_dest(
cls,
db: Session,
dest: str,
storage: Optional[str] = None,
recursive: bool = False,
) -> List["TransferHistory"]:
"""
按目标路径查询成功移动记录,供从媒体库现址发起重新整理时识别历史。
:param db: 数据库会话
:param dest: 目标路径
:param storage: 目标存储类型
:param recursive: 是否递归匹配目录子项
:return: 命中的成功移动记录
"""
normalized_dest = (
Path(str(dest).replace("\\", "/")).as_posix().rstrip("/") or "/"
)
query = db.query(cls).filter(
cls.status.is_(True),
cls.mode.contains("move"),
)
if recursive:
escaped_dest = (
normalized_dest.replace("\\", "\\\\")
.replace("%", "\\%")
.replace("_", "\\_")
)
query = query.filter(
or_(
cls.dest == normalized_dest,
cls.dest.like(f"{escaped_dest.rstrip('/')}/%", escape="\\"),
)
)
else:
query = query.filter(cls.dest == normalized_dest)
if storage:
query = query.filter(cls.dest_storage == storage)
return query.all()
@classmethod
@db_query
def list_by_hash(cls, db: Session, download_hash: str):
return db.query(cls).filter(cls.download_hash == download_hash).all()
@classmethod
@db_query
def statistic(cls, db: Session, days: Optional[int] = 7):
"""
统计最近days天的下载历史数量,按日期分组返回每日数量
"""
sub_query = db.query(func.substr(cls.date, 1, 10).label('date'),
cls.id.label('id')).filter(
cls.date >= time.strftime("%Y-%m-%d %H:%M:%S",
time.localtime(time.time() - 86400 * days))).subquery()
return db.query(sub_query.c.date, func.count(sub_query.c.id)).group_by(sub_query.c.date).all()
@classmethod
@db_query
def monthly_media_statistics(cls, db: Session):
"""
统计当月成功整理的电影、电视剧和剧集数量。
电影和电视剧按媒体身份去重;剧集优先按历史记录中的集数字段计算,
缺少集数时按单条成功整理记录计数。
"""
month_prefix = time.strftime("%Y-%m-", time.localtime())
histories = db.query(cls).filter(
cls.status.is_(True),
cls.date.like(f"{month_prefix}%"),
cls.type.in_([MediaType.MOVIE.value, MediaType.TV.value]),
).all()
movie_identities = set()
tv_identities = set()
episode_count = 0
for history in histories:
identity = (history.tmdbid or 0, history.title or "", history.year or "")
if history.type == MediaType.MOVIE.value:
movie_identities.add(identity)
continue
tv_identities.add(identity)
episode_count += cls._history_episode_count(history)
return len(movie_identities), len(tv_identities), episode_count
@staticmethod
def _history_episode_count(history: "TransferHistory") -> int:
"""从单条整理历史中估算成功入库的剧集数量。"""
episode_numbers = [int(value) for value in re.findall(r"\d+", history.episodes or "")]
if len(episode_numbers) >= 2 and "-" in (history.episodes or ""):
return max(1, episode_numbers[-1] - episode_numbers[0] + 1)
if episode_numbers:
return len(set(episode_numbers))
if isinstance(history.files, list) and history.files:
return len(history.files)
return 1
@classmethod
@async_db_query
async def async_statistic(cls, db: AsyncSession, days: Optional[int] = 7):
"""
统计最近days天的下载历史数量,按日期分组返回每日数量
"""
sub_query = select(func.substr(cls.date, 1, 10).label('date'),
cls.id.label('id')).filter(
cls.date >= time.strftime("%Y-%m-%d %H:%M:%S",
time.localtime(time.time() - 86400 * days))).subquery()
result = await db.execute(
select(sub_query.c.date, func.count(sub_query.c.id)).group_by(sub_query.c.date)
)
return result.all()
@classmethod
@db_query
def count(cls, db: Session, status: bool = None):
if status is not None:
return db.query(func.count(cls.id)).filter(cls.status == status).first()[0]
else:
return db.query(func.count(cls.id)).first()[0]
@classmethod
@async_db_query
async def async_count(cls, db: AsyncSession, status: bool = None):
if status is not None:
result = await db.execute(
select(func.count(cls.id)).filter(cls.status == status)
)
else:
result = await db.execute(
select(func.count(cls.id))
)
return result.scalar()
@classmethod
@db_query
def count_by_title(cls, db: Session, title: str, status: bool = None, wildcard: bool = False):
if wildcard:
text_filter = or_(
_text_like(cls.title, title, wildcard=True),
_text_like(cls.src, title, wildcard=True),
_text_like(cls.dest, title, wildcard=True),
)
else:
text_filter = or_(
_text_like(cls.title, f'%{title}%'),
_text_like(cls.src, f'%{title}%'),
_text_like(cls.dest, f'%{title}%'),
)
query = db.query(func.count(cls.id)).filter(text_filter)
if status is not None:
query = query.filter(cls.status == status)
return query.first()[0]
@classmethod
@async_db_query
async def async_count_by_title(cls, db: AsyncSession, title: str, status: bool = None, wildcard: bool = False):
if wildcard:
text_filter = or_(
_text_like(cls.title, title, wildcard=True),
_text_like(cls.src, title, wildcard=True),
_text_like(cls.dest, title, wildcard=True),
)
else:
text_filter = or_(
_text_like(cls.title, f'%{title}%'),
_text_like(cls.src, f'%{title}%'),
_text_like(cls.dest, f'%{title}%'),
)
stmt = select(func.count(cls.id)).filter(text_filter)
if status is not None:
stmt = stmt.filter(cls.status == status)
result = await db.execute(stmt)
return result.scalar()
@classmethod
@db_query
def list_by(cls, db: Session, mtype: Optional[str] = None, title: Optional[str] = None, year: Optional[str] = None,
season: Optional[str] = None,
episode: Optional[str] = None, tmdbid: Optional[int] = None, dest: Optional[str] = None):
"""
据tmdbid、season、season_episode查询转移记录
tmdbid + mtype 或 title + year 必输
"""
# TMDBID + 类型
if tmdbid and mtype:
# 电视剧某季某集
if season is not None and episode:
return db.query(cls).filter(cls.tmdbid == tmdbid,
cls.type == mtype,
cls.seasons == season,
cls.episodes == episode,
cls.dest == dest).all()
# 电视剧某季
elif season is not None:
return db.query(cls).filter(cls.tmdbid == tmdbid,
cls.type == mtype,
cls.seasons == season).all()
else:
if dest:
# 电影
return db.query(cls).filter(cls.tmdbid == tmdbid,
cls.type == mtype,
cls.dest == dest).all()
else:
# 电视剧所有季集
return db.query(cls).filter(cls.tmdbid == tmdbid,
cls.type == mtype).all()
# 标题 + 年份
elif title and year:
# 电视剧某季某集
if season is not None and episode:
return db.query(cls).filter(cls.title == title,
cls.year == year,
cls.seasons == season,
cls.episodes == episode,
cls.dest == dest).all()
# 电视剧某季
elif season is not None:
return db.query(cls).filter(cls.title == title,
cls.year == year,
cls.seasons == season).all()
else:
if dest:
# 电影
return db.query(cls).filter(cls.title == title,
cls.year == year,
cls.dest == dest).all()
else:
# 电视剧所有季集
return db.query(cls).filter(cls.title == title,
cls.year == year).all()
# 类型 + 转移路径(emby webhook season无tmdbid场景)
elif mtype and season is not None and dest:
# 电视剧某季
return db.query(cls).filter(cls.type == mtype,
cls.seasons == season,
cls.dest.like(f"{dest}%")).all()
return []
@classmethod
@db_query
def get_by_type_tmdbid(cls, db: Session, mtype: Optional[str] = None, tmdbid: Optional[int] = None):
"""
据tmdbid、type查询转移记录
"""
return db.query(cls).filter(cls.tmdbid == tmdbid,
cls.type == mtype).first()
@classmethod
@db_update
def update_download_hash(cls, db: Session, historyid: Optional[int] = None, download_hash: Optional[str] = None):
db.query(cls).filter(cls.id == historyid).update(
{
"download_hash": download_hash
}
)
@classmethod
@db_query
def list_by_date(cls, db: Session, date: str):
"""
查询某时间之后的转移历史
"""
return db.query(cls).filter(cls.date > date).order_by(cls.id.desc()).all()
@classmethod
@db_update
def delete_before(
cls,
db: Session,
before_time: str,
limit: Optional[int] = 500,
) -> int:
"""
分批删除指定时间之前的整理历史。
"""
ids = [
row[0]
for row in db.query(cls.id)
.filter(cls.date < before_time)
.order_by(cls.id.asc())
.limit(limit)
.all()
]
if not ids:
return 0
return (
db.query(cls)
.filter(cls.id.in_(ids))
.delete(synchronize_session=False)
)