| 1 | #!/usr/bin/env python3 |
| 2 | """Import a verified personal Forgejo export into stopped Shale, without replacing native issues.""" |
| 3 | import argparse |
| 4 | from collections import Counter, defaultdict |
| 5 | from datetime import datetime, timezone |
| 6 | import hashlib |
| 7 | import grp |
| 8 | import json |
| 9 | import os |
| 10 | from pathlib import Path |
| 11 | import re |
| 12 | import shutil |
| 13 | import sqlite3 |
| 14 | import sys |
| 15 | import tempfile |
| 16 | import urllib.request |
| 17 | from urllib.parse import quote |
| 18 | |
| 19 | ALPHABET = '0123456789ABCDEFGHJKMNPQRSTVWXYZ' |
| 20 | EPOCH = 1577836800 # Shale's ULIDs use 2020-01-01, rather than the Unix epoch. |
| 21 | CLOSED = ('done', 'not_planned', 'duplicate', 'invalid') |
| 22 | EVENTS = ['comment', 'reopened', 'closed', 'issue reference', 'commit reference', |
| 23 | 'comment reference', 'pull request reference', 'label changed', 'milestone changed', |
| 24 | 'assignee changed', 'title changed', 'branch deleted', 'time tracking started', |
| 25 | 'time tracking stopped', 'time added', 'time tracking canceled', 'deadline added', |
| 26 | 'deadline changed', 'deadline removed', 'dependency added', 'dependency removed', |
| 27 | 'code comment', 'review', 'locked', 'unlocked', 'target branch changed', |
| 28 | 'time deleted', 'review requested', 'merged', 'pull request updated', |
| 29 | 'project changed', 'project column changed', 'review dismissed', 'reference changed', |
| 30 | 'automatic merge scheduled', 'automatic merge canceled', 'pinned', 'unpinned'] |
| 31 | |
| 32 | |
| 33 | def timestamp(epoch): |
| 34 | return datetime.fromtimestamp(int(epoch), timezone.utc).isoformat(timespec='seconds') |
| 35 | |
| 36 | |
| 37 | def identifier(kind, key, epoch): |
| 38 | milliseconds = (int(epoch) - EPOCH) * 1000 |
| 39 | if not 0 <= milliseconds < 2 ** 48: |
| 40 | raise ValueError('Source creation time is outside Shale ULID range') |
| 41 | entropy = int.from_bytes(hashlib.sha256(f'personal-forgejo:{kind}:{key}'.encode()).digest()[:10], 'big') |
| 42 | value = (milliseconds << 80) | entropy |
| 43 | return ''.join(ALPHABET[(value >> (5 * i)) & 31] for i in range(25, -1, -1)) |
| 44 | |
| 45 | |
| 46 | def digest(path): |
| 47 | return hashlib.sha256(path.read_bytes()).hexdigest() |
| 48 | |
| 49 | |
| 50 | def stopped(job): |
| 51 | token = os.environ.get('NOMAD_TOKEN') or Path('/var/lib/studio/nomad.token').read_text().strip() |
| 52 | def read(path): |
| 53 | request = urllib.request.Request('http://127.0.0.1:4646/v1/' + path, |
| 54 | headers={'X-Nomad-Token': token}) |
| 55 | with urllib.request.urlopen(request, timeout=10) as response: |
| 56 | return json.load(response) |
| 57 | if not read('job/' + job).get('Stop') or any( |
| 58 | item['ClientStatus'] in ('pending', 'running') for item in read('job/' + job + '/allocations')): |
| 59 | raise ValueError('Stop the target Shale job and wait for its allocations before importing') |
| 60 | |
| 61 | |
| 62 | def migrate(database, source, attachments, target, evidence, before_commit=lambda: None): |
| 63 | db = sqlite3.connect(database) |
| 64 | db.row_factory = sqlite3.Row |
| 65 | db.execute('PRAGMA foreign_keys=ON') |
| 66 | if db.execute('PRAGMA integrity_check').fetchone()[0] != 'ok': |
| 67 | raise ValueError('Shale SQLite integrity check failed') |
| 68 | source_hash = digest(source) |
| 69 | data = json.loads(source.read_text()) |
| 70 | for field in ['repositories', 'issues', 'comments', 'labels', 'issue_labels', 'authors', |
| 71 | 'attachments', 'assignees', 'milestones', 'pull_requests', 'history']: |
| 72 | if not isinstance(data.get(field), list): |
| 73 | raise ValueError(f'Missing source table: {field}') |
| 74 | owners = db.execute("SELECT * FROM users WHERE name='clover'").fetchall() |
| 75 | if len(owners) != 1 or owners[0]['snowflake'] != '8c909fab-98b0-4472-9171-9606ccebe9d5': |
| 76 | raise ValueError('Clover original identity is missing') |
| 77 | owner = owners[0]['id'] |
| 78 | repos = {r['name']: dict(r) for r in db.execute('SELECT * FROM repositories')} |
| 79 | sources = {r['id']: r for r in data['repositories']} |
| 80 | issue_sources = {r['id']: r for r in data['issues']} |
| 81 | if len(issue_sources) != len(data['issues']): |
| 82 | raise ValueError('Duplicate source issue IDs') |
| 83 | for r in data['repositories']: |
| 84 | if r['name'] not in repos or repos[r['name']]['owner'] != owner: |
| 85 | raise ValueError(f'Expected imported Clover repository: {r["name"]}') |
| 86 | if any(c['issue_id'] not in issue_sources for c in data['comments']): |
| 87 | raise ValueError('Orphan source comment') |
| 88 | with sqlite3.connect(evidence / 'before.db') as backup: |
| 89 | db.backup(backup) |
| 90 | originals = {table: [dict(r) for r in db.execute(f'SELECT * FROM {table}')] |
| 91 | for table in ['issues', 'issue_actions', 'issue_labels', 'issues__labels', 'users']} |
| 92 | db.execute('BEGIN IMMEDIATE') |
| 93 | db.execute('CREATE TABLE IF NOT EXISTS studio_forgejo_records(' |
| 94 | 'kind TEXT NOT NULL, source_id TEXT NOT NULL, target_id INTEGER, ' |
| 95 | 'source_json TEXT NOT NULL, PRIMARY KEY(kind,source_id))') |
| 96 | previous = db.execute("SELECT source_json FROM studio_forgejo_records WHERE kind='export' AND source_id='personal'").fetchone() |
| 97 | if previous and json.loads(previous[0])['sha256'] != source_hash: |
| 98 | raise ValueError('A different export was already imported; review the source before merging') |
| 99 | |
| 100 | def saved(kind, key): |
| 101 | row = db.execute('SELECT target_id FROM studio_forgejo_records WHERE kind=? AND source_id=?', (kind, str(key))).fetchone() |
| 102 | return row[0] if row else None |
| 103 | |
| 104 | def record(kind, key, target_id, row): |
| 105 | db.execute('INSERT OR IGNORE INTO studio_forgejo_records VALUES(?,?,?,?)', |
| 106 | (kind, str(key), target_id, json.dumps(row, ensure_ascii=False, sort_keys=True))) |
| 107 | |
| 108 | users = {1: owner} |
| 109 | names = {r['name'] for r in db.execute('SELECT name FROM users')} |
| 110 | actors = {int(r.get('poster_id') or 0) for r in [*data['issues'], *data['comments'], *data['history']]} |
| 111 | actors.update(int(r['assignee_id']) for r in data['assignees']) |
| 112 | author_source = {r['id']: r for r in data['authors']} |
| 113 | for key in sorted(actors - {1}): |
| 114 | row = author_source.get(key, {'id': key, 'name': 'forgejo-system' if key == 0 else f'forgejo-user-{key}', |
| 115 | 'created_unix': min(i['created_unix'] for i in data['issues']), |
| 116 | 'updated_unix': max(i['updated_unix'] for i in data['issues'])}) |
| 117 | existing = saved('user', key) |
| 118 | if existing is None: |
| 119 | name = row['name'] |
| 120 | if name in names: |
| 121 | name = 'forgejo-' + name |
| 122 | if name in names: |
| 123 | raise ValueError('Historical author name collision') |
| 124 | cursor = db.execute('INSERT INTO users(uuid,provider,snowflake,name,joined_on,last_updated) VALUES(?,?,?,?,?,?)', |
| 125 | (identifier('user', key, row['created_unix']), 'personal-forgejo.invalid', |
| 126 | str(key), name, timestamp(row['created_unix']), timestamp(row['updated_unix']))) |
| 127 | existing = cursor.lastrowid |
| 128 | names.add(name) |
| 129 | record('user', key, existing, row) |
| 130 | users[key] = existing |
| 131 | record('user', 1, owner, author_source[1]) |
| 132 | |
| 133 | external = {'paperclover': owner} |
| 134 | external_rows = defaultdict(list) |
| 135 | for row in [*data['issues'], *data['comments']]: |
| 136 | if row.get('original_author'): |
| 137 | external_rows[row['original_author']].append(row) |
| 138 | for name, rows in external_rows.items(): |
| 139 | if name == 'paperclover': |
| 140 | continue |
| 141 | if not re.fullmatch(r'[A-Za-z0-9_.-]+', name): |
| 142 | raise ValueError('Unsupported external author name') |
| 143 | key = 'external:' + name |
| 144 | existing = saved('user', key) |
| 145 | if existing is None: |
| 146 | display = name if name not in names else 'github-' + name |
| 147 | if display in names: |
| 148 | raise ValueError('External author name collision') |
| 149 | created = min(row['created_unix'] for row in rows) |
| 150 | updated = max(row['updated_unix'] or row['created_unix'] for row in rows) |
| 151 | existing = db.execute('INSERT INTO users(uuid,provider,snowflake,name,joined_on,last_updated) VALUES(?,?,?,?,?,?)', |
| 152 | (identifier('user', key, created), 'personal-forgejo-external.invalid', name, |
| 153 | display, timestamp(created), timestamp(updated))).lastrowid |
| 154 | names.add(display) |
| 155 | record('user', key, existing, {'original_author': name}) |
| 156 | external[name] = existing |
| 157 | |
| 158 | def actor(row): |
| 159 | return external[row['original_author']] if row.get('original_author') else users[row['poster_id']] |
| 160 | |
| 161 | labels = {} |
| 162 | for row in sorted(data['labels'], key=lambda r: r['id']): |
| 163 | if row['repo_id'] not in sources: |
| 164 | raise ValueError('Label has no source repository') |
| 165 | repo = repos[sources[row['repo_id']]['name']] |
| 166 | existing = saved('label', row['id']) |
| 167 | if existing is None: |
| 168 | color = row['color'].lstrip('#') |
| 169 | if not re.fullmatch('[0-9a-fA-F]{6}', color): |
| 170 | raise ValueError('Invalid source label color') |
| 171 | epoch = row['created_unix'] or sources[row['repo_id']]['created_unix'] |
| 172 | existing = db.execute('INSERT INTO issue_labels(uuid,repo,name,description,color,last_updated) VALUES(?,?,?,?,?,?)', |
| 173 | (identifier('label', row['id'], epoch), repo['id'], row['name'], row['description'], |
| 174 | '#' + color, timestamp(row['updated_unix'] or epoch))).lastrowid |
| 175 | record('label', row['id'], existing, row) |
| 176 | labels[row['id']] = existing |
| 177 | |
| 178 | # Preserve incoming numbers when free; collisions go above BOTH namespaces. |
| 179 | occupied = defaultdict(set) |
| 180 | ceilings = defaultdict(int) |
| 181 | for row in db.execute('SELECT repo,number FROM issues'): |
| 182 | occupied[row['repo']].add(row['number']) |
| 183 | ceilings[row['repo']] = max(ceilings[row['repo']], row['number']) |
| 184 | for row in data['issues']: |
| 185 | repo = repos[sources[row['repo_id']]['name']]['id'] |
| 186 | ceilings[repo] = max(ceilings[repo], row['index']) |
| 187 | issue_map = {} |
| 188 | collisions = [] |
| 189 | assignees = defaultdict(list) |
| 190 | for row in data['assignees']: |
| 191 | assignees[row['issue_id']].append(users[row['assignee_id']]) |
| 192 | planned = [] |
| 193 | for row in sorted(data['issues'], key=lambda r: (r['repo_id'], r['index'])): |
| 194 | repo = repos[sources[row['repo_id']]['name']] |
| 195 | existing = saved('issue', row['id']) |
| 196 | number = row['index'] |
| 197 | if existing is None: |
| 198 | if number in occupied[repo['id']]: |
| 199 | ceilings[repo['id']] += 1 |
| 200 | number = ceilings[repo['id']] |
| 201 | occupied[repo['id']].add(number) |
| 202 | else: |
| 203 | number = db.execute('SELECT number FROM issues WHERE id=?', (existing,)).fetchone()[0] |
| 204 | planned.append((repo, number, existing, row)) |
| 205 | # Shale allocates the next number from the issue with the highest row ID. |
| 206 | for repo, number, existing, row in sorted(planned, key=lambda item: (item[0]['id'], item[1])): |
| 207 | if existing is None: |
| 208 | existing = db.execute('INSERT INTO issues(uuid,repo,number,author,last_updated,title,status,assignee) VALUES(?,?,?,?,?,?,?,?)', |
| 209 | (identifier('issue', row['id'], row['created_unix']), repo['id'], number, |
| 210 | actor(row), timestamp(row['updated_unix']), row['name'], |
| 211 | 'done' if row['is_closed'] else 'todo', |
| 212 | assignees[row['id']][0] if len(assignees[row['id']]) == 1 else None)).lastrowid |
| 213 | record('issue', row['id'], existing, row) |
| 214 | if number != row['index']: |
| 215 | collisions.append({'repo': repo['name'], 'forgejo': row['index'], 'shale': number}) |
| 216 | issue_map[row['id']] = {'id': existing, 'number': number, 'repo': repo['name'], |
| 217 | 'url': '/' + quote(repo['name'], safe='/') + '/issues/' + str(number)} |
| 218 | |
| 219 | # Preserve attachment authorization by checking the associated Shale issue |
| 220 | # before Caddy serves the copied file. No unauthenticated static directory. |
| 221 | attachment_map = {} |
| 222 | manifest = [] |
| 223 | for row in data['attachments']: |
| 224 | if row['issue_id'] not in issue_map: |
| 225 | raise ValueError('An attachment belongs to an unsupported release or missing issue') |
| 226 | if row['external_url']: |
| 227 | attachment_map[row['uuid']] = row['external_url'] |
| 228 | record('attachment', row['id'], None, row) |
| 229 | continue |
| 230 | key = row['uuid'] |
| 231 | if not re.fullmatch('[0-9a-f-]{36}', key): |
| 232 | raise ValueError('Unsafe attachment UUID') |
| 233 | source_file = attachments / key[0] / key[1] / key |
| 234 | if source_file.is_symlink() or not source_file.is_file() or source_file.stat().st_size != row['size']: |
| 235 | raise ValueError('Source attachment missing or size differs') |
| 236 | name = Path(row['name']).name |
| 237 | if name != row['name'] or not name or any(c in name for c in '\\"\r\n'): |
| 238 | raise ValueError('Unsafe attachment filename') |
| 239 | destination = target / 'forgejo-attachments' / key / name |
| 240 | destination.parent.mkdir(parents=True, exist_ok=True) |
| 241 | if destination.exists() and digest(destination) != digest(source_file): |
| 242 | raise ValueError('Attachment copy differs') |
| 243 | if not destination.exists(): |
| 244 | shutil.copyfile(source_file, destination) |
| 245 | checksum = digest(source_file) |
| 246 | if digest(destination) != checksum: |
| 247 | raise ValueError('Attachment checksum failed') |
| 248 | # Caddy serves these only after Shale authorizes the associated issue. |
| 249 | # Keep them unavailable to other local users. |
| 250 | caddy_group = grp.getgrnam('caddy').gr_gid |
| 251 | for directory in [target / 'forgejo-attachments', destination.parent]: |
| 252 | os.chown(directory, 0, caddy_group) |
| 253 | directory.chmod(0o750) |
| 254 | os.chown(destination, 0, caddy_group) |
| 255 | destination.chmod(0o640) |
| 256 | path = '/-/forgejo-attachments/' + key + '/' + quote(name, safe='') |
| 257 | attachment_map[key] = path |
| 258 | manifest.append({'path': path, 'file': key + '/' + name, |
| 259 | 'issue': issue_map[row['issue_id']]['url'], 'sha256': checksum}) |
| 260 | record('attachment', row['id'], None, row) |
| 261 | |
| 262 | source_urls = {} |
| 263 | for row in data['issues']: |
| 264 | repo = sources[row['repo_id']] |
| 265 | owner_name = repo['owner_name'] |
| 266 | for kind in ['issues', 'pulls']: |
| 267 | source_urls[f'/{owner_name}/{repo["name"]}/{kind}/{row["index"]}'] = issue_map[row['id']]['url'] |
| 268 | |
| 269 | def rewrite(text, repo_id): |
| 270 | def attachment(match): |
| 271 | key = match.group(1) |
| 272 | if key not in attachment_map: |
| 273 | raise ValueError('Issue references an attachment absent from the source export') |
| 274 | return attachment_map[key] |
| 275 | text = re.sub(r'(?:https?://(?:git|forgejo)\.paperclover\.net)?/attachments/([0-9a-f-]{36})', attachment, text) |
| 276 | def link(match): |
| 277 | return source_urls.get(match.group(1), match.group(0)) |
| 278 | text = re.sub(r'https?://(?:git|forgejo)\.paperclover\.net(/[^\s)<>]+/(?:issues|pulls)/[0-9]+)', link, text) |
| 279 | text = re.sub(r'(?<=\]\()(/[^\s)<>]+/(?:issues|pulls)/[0-9]+)(?=\))', link, text) |
| 280 | # Leave code, quoted text, and bare #references untouched; map explicit links only. |
| 281 | return text |
| 282 | |
| 283 | def action(kind, key, issue_id, actor, payload, added, updated=None, source_row=None): |
| 284 | existing = saved(kind, key) |
| 285 | if existing is not None: |
| 286 | return existing |
| 287 | created = db.execute('INSERT INTO issue_actions(uuid,issue,actor,kind,payload,added_on,last_updated) VALUES(?,?,?,?,?,?,?)', |
| 288 | (identifier(kind, key, added), issue_id, actor, kind if kind in ['comment', 'update_status', 'add_label', 'remove_label'] else 'comment', |
| 289 | payload, timestamp(added), timestamp(updated or added))).lastrowid |
| 290 | record(kind, key, created, source_row or {}) |
| 291 | return created |
| 292 | |
| 293 | milestone_source = {r['id']: r for r in data['milestones']} |
| 294 | pulls = {r['issue_id']: r for r in data['pull_requests']} |
| 295 | body_edits = defaultdict(int) |
| 296 | for revision in data['history']: |
| 297 | if not revision['comment_id'] and not revision['is_first_created']: |
| 298 | body_edits[revision['issue_id']] = max(body_edits[revision['issue_id']], revision['edited_unix']) |
| 299 | for row in data['issues']: |
| 300 | item = issue_map[row['id']] |
| 301 | body = rewrite(row['content'], row['repo_id']) |
| 302 | metadata = [] |
| 303 | if item['number'] != row['index']: |
| 304 | metadata.append(f'Original Forgejo issue #{row["index"]}.') |
| 305 | if row['is_pull']: |
| 306 | pr = pulls[row['id']] |
| 307 | metadata.append(f'Imported pull request: `{pr["head_branch"]}` → `{pr["base_branch"]}`.' + |
| 308 | (f' Merged as `{pr["merged_commit_id"]}`.' if pr['has_merged'] else '')) |
| 309 | if row['milestone_id']: |
| 310 | metadata.append('Milestone: ' + milestone_source[row['milestone_id']]['name'] + '.') |
| 311 | if row['is_locked']: |
| 312 | metadata.append('The original discussion was locked.') |
| 313 | if len(assignees[row['id']]) > 1: |
| 314 | original_assignees = [author_source[a['assignee_id']]['name'] for a in data['assignees'] if a['issue_id'] == row['id']] |
| 315 | metadata.append('Original assignees: ' + ', '.join(original_assignees) + '.') |
| 316 | if metadata: |
| 317 | body += '\n\n---\n\n' + '\n\n'.join(metadata) |
| 318 | action('issue-body', row['id'], item['id'], actor(row), body, |
| 319 | row['created_unix'], max(row['created_unix'], body_edits[row['id']]), row) |
| 320 | # Shale's status-change renderer requires an initial status action. |
| 321 | # Forgejo stores initial state on the issue, rather than as a comment. |
| 322 | action('update_status', 'initial-' + str(row['id']), item['id'], actor(row), |
| 323 | 'todo', row['created_unix']) |
| 324 | |
| 325 | for row in sorted(data['comments'], key=lambda r: (r['created_unix'], r['id'])): |
| 326 | item = issue_map[row['issue_id']] |
| 327 | kind = 'comment' |
| 328 | payload = rewrite(row['content'], issue_sources[row['issue_id']]['repo_id']) |
| 329 | if row['type'] in [1, 2]: |
| 330 | kind, payload = 'update_status', 'todo' if row['type'] == 1 else 'done' |
| 331 | elif row['type'] == 7 and row['label_id'] in labels: |
| 332 | kind = 'add_label' if row['content'] == '1' else 'remove_label' |
| 333 | payload = str(labels[row['label_id']]) |
| 334 | elif row['type'] != 0: |
| 335 | event = EVENTS[row['type']] if row['type'] < len(EVENTS) else f'event {row["type"]}' |
| 336 | detail = [] |
| 337 | for field in ['old_title', 'new_title', 'old_ref', 'new_ref', 'commit_sha', 'commit_id', 'tree_path', 'line']: |
| 338 | if row.get(field): |
| 339 | detail.append(f'{field}: {row[field]}') |
| 340 | for field in ['old_milestone_id', 'milestone_id']: |
| 341 | if row.get(field): |
| 342 | detail.append(field + ': ' + milestone_source.get(row[field], {}).get('name', str(row[field]))) |
| 343 | if row.get('dependent_issue_id'): |
| 344 | linked = issue_map.get(row['dependent_issue_id']) |
| 345 | detail.append('Dependency: ' + (f'[{linked["repo"]}#{linked["number"]}]({linked["url"]})' if linked else str(row['dependent_issue_id']))) |
| 346 | if row.get('ref_issue_id') and row['ref_issue_id'] in issue_map: |
| 347 | linked = issue_map[row['ref_issue_id']] |
| 348 | detail.append(f'[{linked["repo"]}#{linked["number"]}]({linked["url"]})') |
| 349 | payload = 'Forgejo: ' + event + '.' + ('\n\n' + '\n\n'.join(detail) if detail else '') + ('\n\n' + payload if payload else '') |
| 350 | existing = saved('forgejo-comment', row['id']) |
| 351 | if existing is None: |
| 352 | # Every original comment/event maps to exactly one Shale action. |
| 353 | action_id = action(kind, 'forgejo-' + str(row['id']), item['id'], actor(row), payload, |
| 354 | row['created_unix'], row['updated_unix'], row) |
| 355 | record('forgejo-comment', row['id'], action_id, row) |
| 356 | |
| 357 | for row in data['issue_labels']: |
| 358 | item = issue_map[row['issue_id']] |
| 359 | existing = saved('issue-label', row['id']) |
| 360 | if existing is None: |
| 361 | label = labels[row['label_id']] |
| 362 | epoch = issue_sources[row['issue_id']]['updated_unix'] |
| 363 | existing = db.execute('INSERT INTO issues__labels(uuid,issue,label) VALUES(?,?,?)', |
| 364 | (identifier('issue-label', row['id'], epoch), item['id'], label)).lastrowid |
| 365 | record('issue-label', row['id'], existing, row) |
| 366 | |
| 367 | for table in ['history', 'milestones', 'pull_requests', 'assignees']: |
| 368 | for row in data[table]: |
| 369 | record(table, row['id'], None, row) |
| 370 | for source_repo in data['repositories']: |
| 371 | repo = repos[source_repo['name']] |
| 372 | if any(i['repo_id'] == source_repo['id'] for i in data['issues']): |
| 373 | db.execute("UPDATE repositories SET access_issues=? WHERE id=? AND access_issues='off'", |
| 374 | ('private' if source_repo['is_private'] else 'public', repo['id'])) |
| 375 | for field in ['access_issues_submit', 'access_issues_comment']: |
| 376 | db.execute(f"UPDATE repositories SET {field}='private' WHERE id=? AND {field}='off'", (repo['id'],)) |
| 377 | slots = ','.join('?' for _ in CLOSED) |
| 378 | db.execute(f'UPDATE repositories SET open_issues=(SELECT count(*) FROM issues WHERE repo=repositories.id AND status NOT IN ({slots}))', CLOSED) |
| 379 | db.execute(f'UPDATE issue_labels SET open_issues=(SELECT count(*) FROM issues__labels il JOIN issues i ON i.id=il.issue WHERE il.label=issue_labels.id AND i.status NOT IN ({slots}))', CLOSED) |
| 380 | for table, rows in originals.items(): |
| 381 | for row in rows: |
| 382 | current = dict(db.execute(f'SELECT * FROM {table} WHERE id=?', (row['id'],)).fetchone()) |
| 383 | if table == 'issue_labels': |
| 384 | current.pop('open_issues');row = {k:v for k,v in row.items() if k != 'open_issues'} |
| 385 | if current != row: |
| 386 | raise ValueError(f'Existing Shale {table} row changed') |
| 387 | if db.execute('PRAGMA foreign_key_check').fetchall(): |
| 388 | raise ValueError('Imported data has invalid foreign keys') |
| 389 | if db.execute('SELECT 1 FROM issues i WHERE i.id=(SELECT max(id) FROM issues WHERE repo=i.repo) ' |
| 390 | 'AND i.number<>(SELECT max(number) FROM issues WHERE repo=i.repo)').fetchone(): |
| 391 | raise ValueError('Issue row order would reuse a number; repair it before importing') |
| 392 | for kind, rows in [('issue', data['issues']), ('forgejo-comment', data['comments']), ('label', data['labels']), |
| 393 | ('issue-label', data['issue_labels']), ('attachment', data['attachments']), ('history', data['history'])]: |
| 394 | count = db.execute('SELECT count(*) FROM studio_forgejo_records WHERE kind=?', (kind,)).fetchone()[0] |
| 395 | if count != len(rows): |
| 396 | raise ValueError(f'Incomplete import: {kind}') |
| 397 | record('export', 'personal', None, {'sha256': source_hash}) |
| 398 | before_commit() |
| 399 | db.commit() |
| 400 | db.close() |
| 401 | report = {'source_sha256': source_hash, 'issues': len(data['issues']), 'comments_and_events': len(data['comments']), |
| 402 | 'labels': len(labels), 'attachments': len(manifest), 'history_revisions': len(data['history']), |
| 403 | 'number_collisions': collisions, 'issue_mapping': issue_map} |
| 404 | (evidence / 'report.json').write_text(json.dumps(report, indent=2) + '\n') |
| 405 | (target / 'forgejo-attachments.json').write_text(json.dumps(manifest, indent=2) + '\n') |
| 406 | print(json.dumps({k:v for k,v in report.items() if k != 'issue_mapping'}, indent=2)) |
| 407 | |
| 408 | |
| 409 | def repair_issue_order(database, evidence, before_commit): |
| 410 | with sqlite3.connect(database) as db: |
| 411 | db.row_factory = sqlite3.Row |
| 412 | db.execute('PRAGMA foreign_keys=ON') |
| 413 | with sqlite3.connect(evidence / 'before.db') as backup: |
| 414 | db.backup(backup) |
| 415 | db.execute('BEGIN IMMEDIATE') |
| 416 | db.execute('PRAGMA defer_foreign_keys=ON') |
| 417 | if db.execute('PRAGMA integrity_check').fetchone()[0] != 'ok' or db.execute('PRAGMA foreign_key_check').fetchall(): |
| 418 | raise ValueError('Shale database integrity failed') |
| 419 | if db.execute('SELECT 1 FROM issues GROUP BY repo,number HAVING count(*)>1').fetchone(): |
| 420 | raise ValueError('Duplicate issue numbers require separate review') |
| 421 | tables = [r[0] for r in db.execute("SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%'")] |
| 422 | if any(not re.fullmatch('[a-z_]+', table) for table in tables): |
| 423 | raise ValueError('Unexpected table name') |
| 424 | references = {table: [r['from'] for r in db.execute(f'PRAGMA foreign_key_list({table})') if r['table'] == 'issues'] for table in tables} |
| 425 | for table in tables: |
| 426 | for column in db.execute(f'PRAGMA table_info({table})'): |
| 427 | if re.fullmatch('(.*_)?issue(_id)?', column['name']) and column['name'] not in references[table]: |
| 428 | raise ValueError('Undeclared issue reference requires review') |
| 429 | imported = {r[0] for r in db.execute("SELECT target_id FROM studio_forgejo_records WHERE kind='issue'")} |
| 430 | swaps = {} |
| 431 | for repo in db.execute('SELECT DISTINCT repo FROM issues'): |
| 432 | newest = db.execute('SELECT id,number FROM issues WHERE repo=? ORDER BY id DESC LIMIT 1', repo).fetchone() |
| 433 | highest = db.execute('SELECT id,number FROM issues WHERE repo=? ORDER BY number DESC LIMIT 1', repo).fetchone() |
| 434 | if newest['number'] != highest['number']: |
| 435 | if newest['id'] not in imported or highest['id'] not in imported: |
| 436 | raise ValueError('Repair would change a native issue ID') |
| 437 | swaps[newest['id']], swaps[highest['id']] = highest['id'], newest['id'] |
| 438 | |
| 439 | def contents(): |
| 440 | identities = dict(db.execute('SELECT id,uuid FROM issues')) |
| 441 | result = {} |
| 442 | for table in tables: |
| 443 | rows = [] |
| 444 | for record in db.execute(f'SELECT * FROM {table}'): |
| 445 | row = dict(record) |
| 446 | if table == 'issues': |
| 447 | row['id'] = identities[row['id']] if row['id'] in imported else row['id'] |
| 448 | for column in references[table]: |
| 449 | row[column] = identities[row[column]] |
| 450 | if table == 'studio_forgejo_records' and row['kind'] == 'issue': |
| 451 | row['target_id'] = identities[row['target_id']] |
| 452 | rows.append(tuple(row.items())) |
| 453 | result[table] = Counter(rows) |
| 454 | return result |
| 455 | |
| 456 | before = contents() |
| 457 | for old, new in [*((old, -old) for old in swaps), *((-old, new) for old, new in swaps.items())]: |
| 458 | db.execute('UPDATE issues SET id=? WHERE id=?', (new, old)) |
| 459 | for table, columns in references.items(): |
| 460 | for column in columns: |
| 461 | db.execute(f'UPDATE {table} SET {column}=? WHERE {column}=?', (new, old)) |
| 462 | db.execute("UPDATE studio_forgejo_records SET target_id=? WHERE kind='issue' AND target_id=?", (new, old)) |
| 463 | if before != contents() or db.execute('PRAGMA foreign_key_check').fetchall(): |
| 464 | raise ValueError('Repair changed issue content or references') |
| 465 | before_commit() |
| 466 | db.commit() |
| 467 | report = {'surrogate_id_swaps': swaps, 'all_table_contents_preserved': True} |
| 468 | (evidence / 'report.json').write_text(json.dumps(report, indent=2) + '\n') |
| 469 | print(json.dumps(report)) |
| 470 | |
| 471 | |
| 472 | def main(): |
| 473 | parser = argparse.ArgumentParser() |
| 474 | parser.add_argument('--target', type=Path, required=True) |
| 475 | parser.add_argument('--source', type=Path, required=True) |
| 476 | parser.add_argument('--proof', type=Path, help='Verified export proof; defaults to source-proof.json beside the source') |
| 477 | parser.add_argument('--attachments', type=Path, default=Path('/mnt/storage1/apps/forgejo/work/attachments')) |
| 478 | parser.add_argument('--rehearsal', action='store_true') |
| 479 | parser.add_argument('--repair-issue-order', action='store_true') |
| 480 | args = parser.parse_args() |
| 481 | if os.geteuid() != 0: |
| 482 | parser.error('Run as root on Zenith') |
| 483 | os.umask(0o077) |
| 484 | target = args.target.resolve() |
| 485 | if 'evil' in str(target).lower() or target.is_relative_to('/mnt/storage1/apps'): |
| 486 | parser.error('Use a managed Shale target or isolated copy') |
| 487 | if args.rehearsal: |
| 488 | if target.is_relative_to('/srv/prod'): |
| 489 | parser.error('Rehearsal must use an isolated target') |
| 490 | elif target != Path('/srv/prod/shale'): |
| 491 | parser.error('Production imports must target /srv/prod/shale') |
| 492 | guard = lambda: None |
| 493 | if target == Path('/srv/prod/shale'): |
| 494 | guard = lambda: stopped('shale') |
| 495 | elif target.parent == Path('/srv/staging'): |
| 496 | if not re.fullmatch(r'shale-preview-[0-9a-f]{8}', target.name): |
| 497 | parser.error('Expected a managed Shale stage') |
| 498 | guard = lambda: stopped(target.name) |
| 499 | proof = args.proof or args.source.parent / 'source-proof.json' |
| 500 | if not proof.is_file() or json.loads(proof.read_text()).get('sourceSha256') != digest(args.source): |
| 501 | parser.error('Verified export proof is missing or the source checksum differs') |
| 502 | guard() |
| 503 | evidence = Path(tempfile.mkdtemp(prefix='issue-import-', dir=str(args.source.parent))) |
| 504 | if args.repair_issue_order: |
| 505 | repair_issue_order(target / 'data/astheno.shale.db', evidence, guard) |
| 506 | else: |
| 507 | migrate(target / 'data/astheno.shale.db', args.source, args.attachments, target, evidence, guard) |
| 508 | print('Evidence:', evidence) |
| 509 | |
| 510 | |
| 511 | if __name__ == '__main__': |
| 512 | main() |