diff --git a/apps/firehose-service/src/lib/cache-writer.ts b/apps/firehose-service/src/lib/cache-writer.ts index b2acf07..9021178 100644 --- a/apps/firehose-service/src/lib/cache-writer.ts +++ b/apps/firehose-service/src/lib/cache-writer.ts @@ -554,10 +554,12 @@ export async function handleSiteCreateOrUpdate( logger.debug(`Updated site cache for ${did}/${rkey} with record CID ${recordCid}`); // Backfill settings if a record exists for this rkey + // Always skip settings invalidation here - the 'update' invalidation below + // already clears everything including the settings cache on the hosting service const settingsRecord = await fetchSettingsRecord(did, rkey, pdsEndpoint); if (settingsRecord) { await handleSettingsUpdate(did, rkey, settingsRecord.record, settingsRecord.cid, { - skipInvalidation: options?.skipInvalidation, + skipInvalidation: true, }); } diff --git a/apps/hosting-service/src/lib/cache-invalidation.ts b/apps/hosting-service/src/lib/cache-invalidation.ts index a186903..eb6b0ea 100644 --- a/apps/hosting-service/src/lib/cache-invalidation.ts +++ b/apps/hosting-service/src/lib/cache-invalidation.ts @@ -7,13 +7,38 @@ */ import Redis from 'ioredis'; -import { storage } from './storage'; +import type { StorageTier } from '@wispplace/tiered-storage'; +import { hotTier, warmTier } from './storage'; import { cache } from './cache-manager'; const CHANNEL = 'wisp:cache-invalidate'; let subscriber: Redis | null = null; +/** + * Directly invalidate a tier by listing and deleting all keys with the given prefix. + * Each tier is invalidated independently so a failure in one doesn't block the others. + */ +async function invalidateTier( + tier: StorageTier, + tierName: string, + prefix: string, +): Promise { + try { + const keys: string[] = []; + for await (const key of tier.listKeys(prefix)) { + keys.push(key); + } + if (keys.length > 0) { + await tier.deleteMany(keys); + } + return keys.length; + } catch (err) { + console.error(`[CacheInvalidation] Failed to invalidate ${tierName} tier for prefix ${prefix}:`, err); + return 0; + } +} + export function startCacheInvalidationSubscriber(): void { const redisUrl = process.env.REDIS_URL; if (!redisUrl) { @@ -58,10 +83,18 @@ export function startCacheInvalidationSubscriber(): void { console.log(`[CacheInvalidation] Invalidating ${did}/${rkey} (${action})`); - // Clear tiered storage (hot + warm) for this site const prefix = `${did}/${rkey}/`; - const deleted = await storage.invalidate(prefix); - console.log(`[CacheInvalidation] Cleared ${deleted} keys from tiered storage for ${did}/${rkey}`); + + // Invalidate each tier independently - a failure in one tier + // (e.g. S3 listKeys timeout) must NOT prevent hot/warm from being cleared + const hotDeleted = await invalidateTier(hotTier, 'hot', prefix); + const warmDeleted = warmTier + ? await invalidateTier(warmTier, 'warm', prefix) + : 0; + + console.log( + `[CacheInvalidation] Cleared ${hotDeleted} hot + ${warmDeleted} warm keys for ${did}/${rkey}`, + ); // Clear in-memory caches for this site cache.delete('redirectRules', `${did}:${rkey}`); diff --git a/apps/hosting-service/src/lib/storage.ts b/apps/hosting-service/src/lib/storage.ts index 00212e4..77a7285 100644 --- a/apps/hosting-service/src/lib/storage.ts +++ b/apps/hosting-service/src/lib/storage.ts @@ -106,7 +106,13 @@ class ReadOnlyS3Tier implements StorageTier { } async *listKeys(prefix?: string) { - yield* this.tier.listKeys(prefix); + try { + yield* this.tier.listKeys(prefix); + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + console.warn(`[Storage] S3 listKeys error for prefix ${prefix}: ${msg}`); + // Yield nothing on error - don't break invalidation + } } async getStats() { @@ -152,6 +158,109 @@ class ReadOnlyS3Tier implements StorageTier { } } +// Hot tier TTL (seconds) - safety net so stale entries expire even if invalidation fails +const HOT_CACHE_TTL = parseInt(process.env.HOT_CACHE_TTL || '60', 10); // 60s default + +/** + * Wrapper around MemoryStorageTier that enforces a short per-entry TTL. + * This acts as a safety net: even if cache invalidation fails to clear the + * hot tier, stale entries will expire after HOT_CACHE_TTL seconds. + * + * The TieredStorage defaultTTL (14 days) is too long for the hot tier - + * we want stale hot entries to expire quickly and re-fetch from warm/cold. + */ +class TTLMemoryTier implements StorageTier { + public readonly inner: MemoryStorageTier; + private ttlMs: number; + private insertedAt = new Map(); + + constructor(config: { maxSizeBytes: number; maxItems?: number }, ttlSeconds: number) { + this.inner = new MemoryStorageTier(config); + this.ttlMs = ttlSeconds * 1000; + } + + private isStale(key: string): boolean { + const ts = this.insertedAt.get(key); + if (!ts) return false; + return Date.now() - ts > this.ttlMs; + } + + private async evictIfStale(key: string): Promise { + if (this.isStale(key)) { + await this.inner.delete(key); + this.insertedAt.delete(key); + return true; + } + return false; + } + + async get(key: string) { + if (await this.evictIfStale(key)) return null; + return this.inner.get(key); + } + + async getWithMetadata(key: string) { + if (await this.evictIfStale(key)) return null; + return this.inner.getWithMetadata(key); + } + + async getStream(key: string) { + if (await this.evictIfStale(key)) return null; + return this.inner.getStream(key); + } + + async set(key: string, data: Uint8Array, metadata: StorageMetadata) { + this.insertedAt.set(key, Date.now()); + return this.inner.set(key, data, metadata); + } + + async setStream(key: string, stream: NodeJS.ReadableStream, metadata: StorageMetadata) { + this.insertedAt.set(key, Date.now()); + return this.inner.setStream(key, stream, metadata); + } + + async delete(key: string) { + this.insertedAt.delete(key); + return this.inner.delete(key); + } + + async deleteMany(keys: string[]) { + for (const key of keys) this.insertedAt.delete(key); + return this.inner.deleteMany(keys); + } + + async exists(key: string) { + if (await this.evictIfStale(key)) return false; + return this.inner.exists(key); + } + + async *listKeys(prefix?: string) { + yield* this.inner.listKeys(prefix); + } + + async getMetadata(key: string) { + if (await this.evictIfStale(key)) return null; + return this.inner.getMetadata(key); + } + + async setMetadata(key: string, metadata: StorageMetadata) { + return this.inner.setMetadata(key, metadata); + } + + async getStats() { + return this.inner.getStats(); + } + + async clear() { + this.insertedAt.clear(); + return this.inner.clear(); + } +} + +// Exported for direct access during cache invalidation +export let hotTier: TTLMemoryTier; +export let warmTier: StorageTier | undefined; + /** * Initialize tiered storage * Must be called before serving requests @@ -159,7 +268,6 @@ class ReadOnlyS3Tier implements StorageTier { function initializeStorage(): TieredStorage { // Determine cold tier: S3 if configured, otherwise disk acts as cold let coldTier: StorageTier; - let warmTier: StorageTier | undefined; const diskTier = new DiskStorageTier({ directory: CACHE_DIR, @@ -195,13 +303,18 @@ function initializeStorage(): TieredStorage { console.log('[Storage] S3 not configured - using disk-only mode (disk as cold tier)'); } + // Hot tier with short TTL - entries expire quickly so stale data doesn't persist + hotTier = new TTLMemoryTier( + { maxSizeBytes: HOT_CACHE_SIZE, maxItems: HOT_CACHE_COUNT }, + HOT_CACHE_TTL, + ); + + console.log(`[Storage] Hot tier TTL: ${HOT_CACHE_TTL}s`); + const storage = new TieredStorage({ tiers: { - // Hot tier: In-memory LRU for instant serving - hot: new MemoryStorageTier({ - maxSizeBytes: HOT_CACHE_SIZE, - maxItems: HOT_CACHE_COUNT, - }), + // Hot tier: In-memory LRU with short TTL for instant serving + hot: hotTier, // Warm tier: Disk-based cache (only when S3 is configured) warm: warmTier,