Something went wrong. Try again.
Monorepo for wisp.place. A static site hosting service built on top of the AT Protocol. forked from nekomimi.pet/wisp.place-monorepo
Something went wrong. Try again.
11 kB · 366 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367import type { Directory, Entry, File, Record as FsRecord } from '@wispplace/lexicons/types/place/wisp/fs';import type { Record as SubfsRecord } from '@wispplace/lexicons/types/place/wisp/subfs';import { extractBlobCid, resolveDid, getPdsForDid } from '@wispplace/atproto-utils';import { sanitizePath } from '@wispplace/fs-utils';import { existsSync, mkdirSync, writeFileSync, rmSync, renameSync, readFileSync } from 'fs';import { dirname, join } from 'path';import { gunzipSync } from 'zlib';import { createSpinner, formatBytes, pc } from '../lib/progress.ts';import { loadMetadata, saveMetadata, type SiteMetadata } from '../lib/metadata.ts';
const MAX_CONCURRENT_DOWNLOADS = 20;
export interface PullOptions { site: string; path: string;}
async function fetchRecord(pdsEndpoint: string, did: string, collection: string, rkey: string): Promise<any> { const url = `${pdsEndpoint}/xrpc/com.atproto.repo.getRecord?repo=${encodeURIComponent(did)}&collection=${encodeURIComponent(collection)}&rkey=${encodeURIComponent(rkey)}`; const res = await fetch(url); if (!res.ok) { throw new Error(`Failed to fetch record: ${res.status}`); } return res.json();}
function extractSubfsUris(directory: Directory, currentPath: string = ''): Array<{ uri: string; path: string }> { const uris: Array<{ uri: string; path: string }> = [];
for (const entry of directory.entries) { const fullPath = currentPath ? `${currentPath}/${entry.name}` : entry.name;
if ('type' in entry.node) { if (entry.node.type === 'subfs') { const subfsNode = entry.node as any; if (subfsNode.subject) { uris.push({ uri: subfsNode.subject, path: fullPath }); } } else if (entry.node.type === 'directory') { const subUris = extractSubfsUris(entry.node as Directory, fullPath); uris.push(...subUris); } } }
return uris;}
async function expandSubfsNodes( directory: Directory, pdsEndpoint: string, depth: number = 0, subfsCache: Map<string, SubfsRecord | null> = new Map()): Promise<Directory> { const MAX_DEPTH = 10;
if (depth >= MAX_DEPTH) { console.warn('Max subfs expansion depth reached'); return directory; }
const subfsUris = extractSubfsUris(directory); if (subfsUris.length === 0) { return directory; }
// Fetch uncached subfs records const uncachedUris = subfsUris.filter(({ uri }) => !subfsCache.has(uri));
if (uncachedUris.length > 0) { await Promise.all(uncachedUris.map(async ({ uri }) => { try { const parts = uri.replace('at://', '').split('/'); const did = parts[0]!; const collection = parts[1]!; const rkey = parts[2]!;
const data = await fetchRecord(pdsEndpoint, did, collection, rkey); subfsCache.set(uri, data.value as SubfsRecord); } catch { subfsCache.set(uri, null); } })); }
// Build map of path -> entries const subfsMap = new Map<string, Entry[]>(); for (const { uri, path } of subfsUris) { const record = subfsCache.get(uri); if (record?.root?.entries) { subfsMap.set(path, record.root.entries as unknown as Entry[]); } }
// Replace subfs nodes with their content function replaceSubfsInEntries(entries: Entry[], currentPath: string = ''): Entry[] { const result: Entry[] = [];
for (const entry of entries) { const fullPath = currentPath ? `${currentPath}/${entry.name}` : entry.name; const node = entry.node;
if ('type' in node && node.type === 'subfs') { const subfsNode = node as any; const isFlat = subfsNode.flat !== false; const subfsEntries = subfsMap.get(fullPath);
if (subfsEntries) { if (isFlat) { const processedEntries = replaceSubfsInEntries(subfsEntries, currentPath); result.push(...processedEntries); } else { const processedEntries = replaceSubfsInEntries(subfsEntries, fullPath); result.push({ name: entry.name, node: { type: 'directory', entries: processedEntries } as any }); } } else { result.push(entry); } } else if ('type' in node && node.type === 'directory' && 'entries' in node) { result.push({ ...entry, node: { ...node, entries: replaceSubfsInEntries(node.entries, fullPath) } }); } else { result.push(entry); } }
return result; }
const partiallyExpanded = { ...directory, entries: replaceSubfsInEntries(directory.entries) };
return expandSubfsNodes(partiallyExpanded, pdsEndpoint, depth + 1, subfsCache);}
interface FileToDownload { path: string; cid: string; encoding?: 'gzip'; mimeType?: string; base64?: boolean;}
function collectFiles( entries: Entry[], pathPrefix: string, existingCids: Record<string, string>): { toDownload: FileToDownload[]; toSkip: number } { const toDownload: FileToDownload[] = []; let toSkip = 0;
function collect(entries: Entry[], currentPath: string) { for (const entry of entries) { const fullPath = currentPath ? `${currentPath}/${entry.name}` : entry.name; const node = entry.node;
if ('type' in node && node.type === 'directory' && 'entries' in node) { collect(node.entries, fullPath); } else if ('type' in node && node.type === 'file' && 'blob' in node) { const fileNode = node as File; const cid = extractBlobCid(fileNode.blob);
if (!cid) continue;
if (existingCids[fullPath] === cid) { toSkip++; } else { toDownload.push({ path: fullPath, cid, encoding: fileNode.encoding, mimeType: fileNode.mimeType, base64: fileNode.base64 }); } } } }
collect(entries, pathPrefix); return { toDownload, toSkip };}
async function downloadBlob( pdsEndpoint: string, did: string, file: FileToDownload): Promise<Buffer> { const url = `${pdsEndpoint}/xrpc/com.atproto.sync.getBlob?did=${encodeURIComponent(did)}&cid=${encodeURIComponent(file.cid)}`; const res = await fetch(url);
if (!res.ok) { throw new Error(`Failed to download blob ${file.cid}: ${res.status}`); }
let content = Buffer.from(await res.arrayBuffer());
// Decode base64 if needed if (file.base64) { const base64String = content.toString('utf-8'); content = Buffer.from(base64String, 'base64'); }
// Decompress gzip if (file.encoding === 'gzip' && content.length >= 2 && content[0] === 0x1f && content[1] === 0x8b) { try { content = gunzipSync(content); } catch { // Keep original content if decompression fails } }
return content;}
export async function pull( identifier: string, options: PullOptions): Promise<void> { const { site, path: outputPath } = options;
console.log(pc.cyan(`\nPulling ${pc.bold(site)} from ${identifier}\n`));
// 1. Resolve DID const spinner = createSpinner('Resolving identity...').start(); const did = await resolveDid(identifier);
if (!did) { spinner.fail('Failed to resolve identity'); throw new Error(`Could not resolve: ${identifier}`); }
spinner.succeed(`Resolved to ${did}`);
// 2. Get PDS endpoint const pdsSpinner = createSpinner('Getting PDS endpoint...').start(); const pdsEndpoint = await getPdsForDid(did);
if (!pdsEndpoint) { pdsSpinner.fail('Failed to get PDS endpoint'); throw new Error(`Could not get PDS for: ${did}`); }
pdsSpinner.succeed(`PDS: ${pdsEndpoint}`);
// 3. Fetch site record const recordSpinner = createSpinner('Fetching site record...').start(); let recordData;
try { recordData = await fetchRecord(pdsEndpoint, did, 'place.wisp.fs', site); } catch { recordSpinner.fail('Site not found'); throw new Error(`Site not found: ${site}`); }
const record = recordData.value as FsRecord; const recordCid = recordData.cid || ''; recordSpinner.succeed('Fetched site record');
// 4. Expand subfs nodes const expandSpinner = createSpinner('Expanding subfs nodes...').start(); const expandedRoot = await expandSubfsNodes(record.root, pdsEndpoint); expandSpinner.succeed('Expanded subfs nodes');
// 5. Load existing metadata for incremental updates const existingMetadata = loadMetadata(outputPath); const existingCids = existingMetadata?.fileCids || {};
// 6. Collect files to download const { toDownload, toSkip } = collectFiles(expandedRoot.entries, '', existingCids);
console.log(pc.dim(`Files to download: ${toDownload.length}, unchanged: ${toSkip}`));
if (toDownload.length === 0 && toSkip > 0) { console.log(pc.green('\n✓ Site is already up to date\n')); return; }
// 7. Create temp directory const tempDir = `${outputPath}.tmp-${Date.now()}`; mkdirSync(tempDir, { recursive: true });
// 8. Download files const downloadSpinner = createSpinner(`Downloading ${toDownload.length} files...`).start(); const newFileCids: Record<string, string> = { ...existingCids }; let downloaded = 0;
try { for (let i = 0; i < toDownload.length; i += MAX_CONCURRENT_DOWNLOADS) { const batch = toDownload.slice(i, i + MAX_CONCURRENT_DOWNLOADS);
await Promise.all(batch.map(async (file) => { const content = await downloadBlob(pdsEndpoint, did, file); const filePath = join(tempDir, sanitizePath(file.path));
mkdirSync(dirname(filePath), { recursive: true }); writeFileSync(filePath, content);
newFileCids[file.path] = file.cid; downloaded++; downloadSpinner.text = `Downloading files: ${downloaded}/${toDownload.length}`; })); }
downloadSpinner.succeed(`Downloaded ${downloaded} files`);
// 9. Copy unchanged files from existing directory if (toSkip > 0 && existsSync(outputPath)) { const copySpinner = createSpinner(`Copying ${toSkip} unchanged files...`).start();
for (const [filePath, cid] of Object.entries(existingCids)) { if (!toDownload.find(f => f.path === filePath)) { const srcPath = join(outputPath, sanitizePath(filePath)); const destPath = join(tempDir, sanitizePath(filePath));
if (existsSync(srcPath)) { mkdirSync(dirname(destPath), { recursive: true }); const content = readFileSync(srcPath); writeFileSync(destPath, content); } } }
copySpinner.succeed(`Copied ${toSkip} unchanged files`); }
// 10. Atomic replace if (existsSync(outputPath)) { const backupPath = `${outputPath}.backup-${Date.now()}`; renameSync(outputPath, backupPath); renameSync(tempDir, outputPath); rmSync(backupPath, { recursive: true, force: true }); } else { renameSync(tempDir, outputPath); }
// 11. Save metadata const metadata: SiteMetadata = { recordCid, fileCids: newFileCids, lastSync: Date.now() }; saveMetadata(outputPath, metadata);
console.log(pc.green(`\n✓ Pulled ${site} to ${outputPath}\n`));
} catch (err) { // Cleanup temp dir on error rmSync(tempDir, { recursive: true, force: true }); throw err; }}