Files
ucvl-home-vision/journal.py
T

299 lines
20 KiB
Python
Raw Normal View History

"""Personal journals, household sharing and owner-bound cloud attachments."""
import datetime as dt
import hashlib
import os
import re
import secrets
import threading
import time
CHUNK = 512 * 1024
MAX_FILE = 256 * 1024 * 1024
KINDS = {'日常', '成长', '家庭', '旅行', '工作', '其他'}
TYPES = {'.jpg': 'image/jpeg', '.jpeg': 'image/jpeg', '.png': 'image/png', '.webp': 'image/webp',
'.mp4': 'video/mp4', '.mov': 'video/quicktime', '.webm': 'video/webm',
'.mp3': 'audio/mpeg', '.m4a': 'audio/mp4', '.aac': 'audio/aac', '.wav': 'audio/wav', '.ogg': 'audio/ogg'}
class Journal:
def __init__(self, app):
self.a = app
self.guard = threading.RLock()
self.buffers = {}
def initialize(self):
with self.a.LOCK, self.a.DB:
for table in ('journal_entries', 'journal_media', 'journal_comments', 'journal_likes'):
self.a.DB.execute('CREATE TABLE IF NOT EXISTS '+table+' (id TEXT PRIMARY KEY, body TEXT NOT NULL)')
# No binary files are written to disk. Partial 4 MiB blocks vanish on restart.
def owned(self, user, table, key):
item = self.a.HOUSEHOLDS.owns(user, self.a.get_object(table, key))
if item['ownerId'] != user['id']:
raise self.a.Problem('只能修改自己的记录', 403)
return item
def visible(self, user, item):
return bool(item and not item.get('deleted') and item['familyId'] == self.a.HOUSEHOLDS.scope(user)
and (item['ownerId'] == user['id'] or item['visibility'] == 'household'))
def entry(self, user, key):
item = self.a.get_object('journal_entries', key)
if not self.visible(user, item): raise self.a.Problem('记事不存在或未获授权', 404)
return item
def public_media(self, item):
return {k: item.get(k) for k in ('id', 'name', 'mime', 'size', 'status', 'received', 'progress', 'error', 'at', 'trashedAt')}
def media(self, user, key):
item = self.a.HOUSEHOLDS.owns(user, self.a.get_object('journal_media', key))
if item['ownerId'] != user['id'] and not any(self.visible(user, entry) and key in entry['mediaIds'] for entry in self.a.objects('journal_entries')):
raise self.a.Problem('附件不存在或未获授权', 404)
if item['status'] != 'ready': raise self.a.Problem('附件尚未上传完成', 409)
return item
def snapshot(self, user, query):
a = self.a
if query.get('mode', 'feed') not in ('feed', 'mine', 'person'):
raise a.Problem('查询方式不正确')
entries = [e for e in a.objects('journal_entries') if self.visible(user, e)]
mode = query.get('mode', 'feed')
if mode == 'feed': entries = [e for e in entries if e['visibility'] == 'household']
elif mode == 'mine': entries = [e for e in entries if e['ownerId'] == user['id']]
elif mode == 'person': entries = [e for e in entries if e['ownerId'] == query.get('person')]
for field, compare in [('from', lambda x, y: x >= y), ('to', lambda x, y: x <= y)]:
if query.get(field):
self.date(query[field]); entries = [e for e in entries if compare(e['occurredOn'], query[field])]
search = query.get('search', '').strip().casefold()[:120]
if search: entries = [e for e in entries if search in (e['text']+' '+e['category']).casefold()]
if query.get('category'): entries = [e for e in entries if e['category'] == query['category']]
entries.sort(key=lambda e: (e['occurredOn'], e['createdAt'], e['id']), reverse=True)
try: page = max(0, min(100000, int(query.get('page', '0'))))
except ValueError: raise a.Problem('页码不正确')
if query.get('focus'):
index = next((i for i, entry in enumerate(entries) if entry['id'] == query['focus']), None)
if index is not None: page = index // 20
page = min(page, max(0, (len(entries)-1)//20))
media = {m['id']: m for m in a.HOUSEHOLDS.objects(user, 'journal_media')}
likes = a.HOUSEHOLDS.objects(user, 'journal_likes')
comments = a.HOUSEHOLDS.objects(user, 'journal_comments')
items = []
for entry in entries[page*20:page*20+20]:
related = sorted([c for c in comments if c['entryId'] == entry['id'] and not c.get('deleted')], key=lambda c: c['createdAt'])
voters = [v for v in likes if v['entryId'] == entry['id'] and v['liked']]
items.append(dict(entry, media=[self.public_media(media[mid]) for mid in entry['mediaIds'] if mid in media],
canEdit=entry['ownerId'] == user['id'], likes=len(voters), liked=any(v['ownerId'] == user['id'] for v in voters),
comments=[dict(c, canDelete=c['ownerId'] == user['id'] or entry['ownerId'] == user['id']) for c in related[-100:]],
commentCount=len(related)))
owned = [m for m in media.values() if m['ownerId'] == user['id']]
uploads = [self.public_media(m) for m in owned if not m.get('trashedAt') and m['status'] != 'removed']
uploads.sort(key=lambda m: m['at'], reverse=True)
trash = [self.public_media(m) for m in owned if m.get('trashedAt') or m['status'] == 'removed']
trash.sort(key=lambda m: m.get('trashedAt') or m['at'], reverse=True)
people = {e['ownerId']: e['authorName'] for e in a.objects('journal_entries') if self.visible(user, e)}
return dict(items=items, total=len(entries), page=page, uploads=uploads[:50], trash=trash[:50], people=people,
connection=a.BAIDU.status(user), chunkSize=CHUNK, maxFile=MAX_FILE)
def date(self, value):
try:
if dt.date.fromisoformat(value).isoformat() != value: raise ValueError()
except (TypeError, ValueError): raise self.a.Problem('请填写有效的公历日期')
return value
def save(self, user, data):
a = self.a
with a.LOCK:
old = self.owned(user, 'journal_entries', data['id']) if data.get('id') else None
if old and old.get('deleted'): raise a.Problem('记录已经删除', 404)
if type(data.get('revision', 0)) is not int or data.get('revision', 0) != (old or {}).get('revision', 0):
raise a.Problem('记事已更新,请重新打开后保存', 409)
client = a.clean_text(data.get('clientId', ''), 32, True)
if not re.fullmatch('[a-f0-9]{32}', client): raise a.Problem('提交编号不正确')
if not old:
prior = next((e for e in a.objects('journal_entries') if e['ownerId'] == user['id'] and e['familyId'] == a.HOUSEHOLDS.scope(user) and e['clientId'] == client), None)
if prior: return {'id': prior['id'], 'duplicate': True}
raw_text = data.get('text', '')
if not isinstance(raw_text, str): raise a.Problem('记事内容须为文字')
text = a.clean_text(raw_text.replace('\r\n', '\n').replace('\n', '\u2028'), 5000).replace('\u2028', '\n')
visibility = data.get('visibility', 'private'); category = data.get('category', '日常')
if visibility not in ('private', 'household') or category not in KINDS:
raise a.Problem('可见范围或分类不正确')
ids = data.get('mediaIds', [])
if not isinstance(ids, list) or len(ids) > 9 or any(not isinstance(i, str) for i in ids) or len(set(ids)) != len(ids):
raise a.Problem('每条记事最多 9 个附件')
for mid in ids:
item = self.owned(user, 'journal_media', mid)
if item['status'] != 'ready': raise a.Problem('请等待附件上传完成后再发布')
if item.get('trashedAt') and mid not in (old or {}).get('mediaIds', []):
raise a.Problem('请先恢复已删除的附件再添加到记事', 409)
if not text and not ids: raise a.Problem('请写下内容或添加附件')
now = dt.datetime.now(dt.timezone.utc).isoformat()
entry = dict(id=old['id'] if old else secrets.token_hex(8), familyId=a.HOUSEHOLDS.scope(user),
ownerId=user['id'], authorName=user['name'], text=text, category=category,
visibility=visibility, occurredOn=self.date(data.get('occurredOn')),
mediaIds=ids, createdAt=old['createdAt'] if old else now, updatedAt=now,
revision=(old or {}).get('revision', 0)+1, clientId=client, deleted=False)
a.save_object('journal_entries', entry)
return {'id': entry['id']}
def action(self, user, data):
a = self.a
with a.LOCK:
entry = self.entry(user, data.get('id'))
action = data.get('action')
if action == 'delete':
self.owned(user, 'journal_entries', entry['id'])
if data.get('revision') != entry['revision']: raise a.Problem('记事已更新,请重新读取', 409)
entry.update(deleted=True, revision=entry['revision']+1); a.save_object('journal_entries', entry)
elif action == 'like':
if not isinstance(data.get('liked'), bool): raise a.Problem('点赞状态不正确')
key = hashlib.sha256((entry['id']+user['id']).encode()).hexdigest()[:32]
a.save_object('journal_likes', dict(id=key, entryId=entry['id'], familyId=entry['familyId'], ownerId=user['id'], liked=data['liked']))
elif action == 'comment':
text = a.clean_text(data.get('text', ''), 500, True)
key = a.clean_text(data.get('clientId', ''), 32, True)
if not re.fullmatch('[a-f0-9]{32}', key): raise a.Problem('评论编号不正确')
old = a.get_object('journal_comments', key)
if old:
if old['ownerId'] != user['id'] or old['entryId'] != entry['id']: raise a.Problem('评论编号已使用', 409)
return {'ok': True}
count = sum(c['entryId'] == entry['id'] and not c.get('deleted') for c in a.objects('journal_comments'))
if count >= 100: raise a.Problem('本条记事已达到 100 条评论上限')
a.save_object('journal_comments', dict(id=key, entryId=entry['id'], familyId=entry['familyId'], ownerId=user['id'], authorName=user['name'], text=text, createdAt=dt.datetime.now(dt.timezone.utc).isoformat()))
elif action == 'delete-comment':
comment = a.get_object('journal_comments', data.get('commentId'))
if not comment or comment['entryId'] != entry['id']: raise a.Problem('评论不存在', 404)
if user['id'] not in (entry['ownerId'], comment['ownerId']): raise a.Problem('只能删除自己的评论或自己记事下的评论', 403)
comment['deleted'] = True; a.save_object('journal_comments', comment)
else: raise a.Problem('记事操作不正确')
return {'ok': True}
def begin_upload(self, user, data):
a = self.a
name = a.clean_text(data.get('name', ''), 120, True)
ext = os.path.splitext(name)[1].lower()
if ext not in TYPES or any(c in name for c in '/\\:*?"<>|') or name.startswith('.'):
raise a.Problem('请选择 JPG、PNG、WebP 照片,MP4、MOV、WebM 视频或常见音频文件')
size = a.integer(data.get('size'), 1, MAX_FILE)
hashes = data.get('hashes')
if not isinstance(hashes, list) or len(hashes) != (size+4194303)//4194304 or any(not isinstance(h, str) or not re.fullmatch('[a-f0-9]{32}', h) for h in hashes):
raise a.Problem('文件分片校验表不正确')
mime = 'audio/webm' if ext == '.webm' and data.get('mime', '').startswith('audio/') else TYPES[ext]
a.BAIDU.token('shared')
with self.guard:
for key, buffer in list(self.buffers.items()):
if buffer['touched'] < time.time()-600: self.buffers.pop(key)
prior = self.owned(user, 'journal_media', data['id']) if data.get('id') else None
if prior:
if prior.get('trashedAt'): raise a.Problem('请先恢复已删除的附件', 409)
if prior['hashes'] != hashes or prior['size'] != size or prior['name'] != name:
raise a.Problem('请选择与原任务相同的文件')
if prior['status'] == 'ready': return self.public_media(prior)
if prior['status'] == 'removed': raise a.Problem('该附件已移除,请新建上传任务')
item = prior
else:
active = [m for m in a.objects('journal_media') if m['status'] == 'receiving' and not m.get('trashedAt') and m['ownerId'] == user['id']]
if len(active) >= 4: raise a.Problem('请先完成或移除现有上传任务,最多同时保留 4 个', 409)
key = secrets.token_hex(8); family = a.HOUSEHOLDS.scope(user)
instance = a.setting('cloud_instance')
if not instance:
instance = secrets.token_hex(8); a.set_setting('cloud_instance', instance)
category = '照片' if mime.startswith('image/') else '视频' if mime.startswith('video/') else '语音'
remote = '/apps/'+a.BAIDU.public_config()['appFolder']+'/个人记事/'+instance+'/'+family+'/'+user['id']+'/'+category+'/'+key+'/'+name
item = dict(id=key, familyId=family, ownerId=user['id'], name=name, mime=mime, size=size,
status='receiving', received=0, progress=0, error='', remotePath=remote, hashes=hashes,
at=dt.datetime.now(dt.timezone.utc).isoformat())
a.save_object('journal_media', item)
try:
pre = a.BAIDU.prepare(item['remotePath'], size, hashes)
if pre.get('existing'):
item.update(status='ready', received=size, progress=100, fsId=str(pre['existing']['fs_id']), error='')
else:
# Precreate reports which cloud blocks are still missing; after browser/server
# restart resend from the first missing block, without any local file cache.
needed = pre.get('block_list', list(range(len(hashes))))
if not isinstance(needed, list) or any(type(n) is not int or not 0 <= n < len(hashes) for n in needed):
raise a.Problem('网盘返回的分片索引不正确', 502)
received = min(needed)*4194304 if needed else size
item.update(uploadId=pre['uploadid'], received=received, progress=round(received/size*95), error='')
self.buffers.pop(item['id'], None)
a.save_object('journal_media', item)
except a.Problem as error:
item['error'] = str(error); a.save_object('journal_media', item); raise
return self.public_media(item)
@staticmethod
def valid_signature(mime, raw):
return bool(mime == 'image/jpeg' and raw.startswith(b'\xff\xd8\xff') or
mime == 'image/png' and raw.startswith(b'\x89PNG\r\n\x1a\n') or
mime == 'image/webp' and raw.startswith(b'RIFF') and raw[8:12] == b'WEBP' or
mime in ('video/mp4', 'video/quicktime', 'audio/mp4') and raw[4:8] == b'ftyp' or
mime in ('video/webm', 'audio/webm') and raw.startswith(b'\x1a\x45\xdf\xa3') or
mime == 'audio/wav' and raw.startswith(b'RIFF') and raw[8:12] == b'WAVE' or
mime == 'audio/ogg' and raw.startswith(b'OggS') or
mime in ('audio/mpeg', 'audio/aac') and (raw.startswith(b'ID3') or len(raw) > 1 and raw[0] == 255 and raw[1] & 224 == 224))
def chunk(self, user, key, offset, content):
if not content or len(content) > CHUNK: raise self.a.Problem('文件分片大小不正确')
# At most two 4 MiB buffers; this lock covers only cloud upload, never the app DB.
with self.guard:
item = self.owned(user, 'journal_media', key)
if item.get('trashedAt'): raise self.a.Problem('附件已删除,请先恢复', 409)
if item['status'] != 'receiving' or not item.get('uploadId'): raise self.a.Problem('请重新选择文件继续上传', 409)
if offset < 0 or offset+len(content) > item['size']: raise self.a.Problem('文件分片位置不正确')
if offset < item['received']: return dict(self.public_media(item), received=item['received'])
for mid, value in list(self.buffers.items()):
if value['touched'] < time.time()-600: self.buffers.pop(mid)
if key not in self.buffers:
if len(self.buffers) >= 2: raise self.a.Problem('上传缓冲正在使用,请稍后重试', 429)
self.buffers[key] = dict(start=item['received'], data=bytearray(), touched=time.time())
buffer = self.buffers[key]; end = buffer['start']+len(buffer['data'])
if offset == end: buffer['data'].extend(content)
elif offset < end and buffer['data'][offset-buffer['start']:offset-buffer['start']+len(content)] == content: pass
else: raise self.a.Problem('分片顺序或内容不正确,请重新选择文件', 409)
buffer['touched'] = time.time(); end = buffer['start']+len(buffer['data'])
if len(buffer['data']) > 4194304: raise self.a.Problem('分片超出缓冲上限')
if buffer['start'] == 0 and not self.valid_signature(item['mime'], buffer['data'][:32]):
self.buffers.pop(key); raise self.a.Problem('文件内容与格式不符,请选择原始媒体文件')
if len(buffer['data']) == 4194304 or end == item['size']:
index = buffer['start']//4194304
self.a.BAIDU.send_block(item['remotePath'], item['uploadId'], index, bytes(buffer['data']), item['hashes'][index])
item.update(received=end, progress=round(end/item['size']*95), error='')
self.a.save_object('journal_media', item); self.buffers.pop(key)
return dict(self.public_media(item), received=end)
def finish_upload(self, user, key):
with self.guard:
item = self.owned(user, 'journal_media', key)
if item.get('trashedAt'): raise self.a.Problem('附件已删除,请先恢复', 409)
if item['status'] == 'ready': return self.public_media(item)
if item['status'] != 'receiving' or item['received'] != item['size']:
raise self.a.Problem('文件尚未完整上传,请重新选择原文件继续', 409)
result = self.a.BAIDU.complete(item)
item.update(status='ready', progress=100, fsId=str(result['fs_id']), error='')
self.a.save_object('journal_media', item)
return self.public_media(item)
def remove_upload(self, user, key):
with self.guard, self.a.LOCK:
item = self.owned(user, 'journal_media', key)
self.buffers.pop(key, None)
# Library deletion must not break existing published references or delete cloud files.
if not item.get('trashedAt'):
item['trashedAt'] = dt.datetime.now(dt.timezone.utc).isoformat()
self.a.save_object('journal_media', item)
return {'ok': True}
def restore_upload(self, user, key):
with self.guard, self.a.LOCK:
item = self.owned(user, 'journal_media', key)
if not item.get('trashedAt') and item['status'] != 'removed': return {'ok': True}
status = ('ready' if item.get('fsId') else 'receiving') if item['status'] == 'removed' else item['status']
if status == 'receiving':
active = sum(m['ownerId'] == user['id'] and m['status'] == 'receiving' and not m.get('trashedAt') for m in self.a.objects('journal_media'))
if active >= 4: raise self.a.Problem('请先完成或删除现有上传任务,最多同时保留 4 个', 409)
item.update(status=status, trashedAt=None)
self.a.save_object('journal_media', item)
return {'ok': True}