273 lines
18 KiB
Python
273 lines
18 KiB
Python
"""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')}
|
|
|
|
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('页码不正确')
|
|
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)))
|
|
uploads = [self.public_media(m) for m in media.values() if m['ownerId'] == user['id'] and m['status'] != 'removed']
|
|
uploads.sort(key=lambda m: 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], 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 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['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 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)
|
|
remote = '/apps/'+a.BAIDU.public_config()['appFolder']+'/资料/'+instance+'/'+family+'/'+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['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['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)
|
|
if any(not e.get('deleted') and key in e['mediaIds'] for e in self.a.objects('journal_entries')):
|
|
raise self.a.Problem('附件仍被记事引用,请先编辑记事移除附件', 409)
|
|
self.buffers.pop(key, None)
|
|
item.update(status='removed'); self.a.save_object('journal_media', item)
|
|
return {'ok': True}
|