"""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) 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.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}