diff --git a/docs/TECHNICAL.md b/docs/TECHNICAL.md index 7876fbf..d946691 100644 --- a/docs/TECHNICAL.md +++ b/docs/TECHNICAL.md @@ -71,6 +71,7 @@ flowchart TB | `lib/classic-liq.ts` | Classic asset trustline utilities (`changeTrust` XDR, Horizon checks) + CID integrity helpers via `cid-cache` (phase-119) + rate-limit-aware batch submission to Horizon (phase-134) | | `lib/phase-copy.ts` | Centralized i18n dictionary (EN/ES) | | `lib/server-data-paths.ts` | Writable data location abstraction | +| `lib/json-store.ts` | Serialized read-modify-write + atomic write for the JSON sidecar stores (§5.4) | | `lib/feature-flags.ts` | Flag registry (phase-107,111,113,114 + 116,117,119,120 + 121..124, env resolution, rollback notes) | | `lib/world-conflict.ts` | World save conflict detection (phase-105): version check + per-author vector clock normalization (object / `Map` / entries array) and ordering | | `lib/story-arc-continuity.ts` | AI story-arc continuity check against recent world narratives (phase-107) | @@ -228,6 +229,39 @@ edit takes a pre-edit snapshot into `signal_versions` and must stay revertible. A merge-based CRDT would fold concurrent edits together and destroy that history, so the contract is compare-and-swap rather than automatic merge. +### 5.4 Shared JSON sidecar store layer + +`lib/json-store.ts` provides the same three guarantees described in §5.3 — +serialized mutation, atomic replacement, and fatal corruption — as a reusable +module, and the remaining JSON sidecar stores route their mutations through it +rather than hand-rolling a `readFile`/`writeFile` pair per module. + +| Export | Purpose | +|---|---| +| `updateStore(key, mutate)` | Locked read-modify-write for a registered sidecar file. The normal way to mutate one. | +| `updateStoreWithReader(key, read, mutate)` | As above, with an on-disk shape adapter applied before mutation. | +| `updateJsonFile(path, { read?, mutate })` | As above, for a path outside the `serverDataJsonPath` registry. | +| `readStore(key, fallback)` / `readJsonFile(path, fallback)` | Unlocked reads. Missing file → `fallback`; corrupt file → throws. | +| `writeJsonFileAtomic(path, data)` | Unlocked atomic write. Torn-write safety only, not lost-update safety. | +| `withFileLock(path, task)` | Escape hatch. Almost always prefer `updateStore`, which keeps the read and write in one critical section. | + +Stores migrated: `follow-store`, `profile-store`, `notification-store`, +`achievement-store`, `narrative-world-store`, and the JSON-backed paths in +`market-store` (`marketProfileViews`, `blockList`; listings and offers are +already SQLite). + +`updateStore` accepts an async `mutate`, but **the lock is not reentrant**: a +`mutate` callback must not call a locking store function for the same file, or it +deadlocks waiting on its own queue. Two places depend on this contract — +`createNotificationBatch` resolves notification preferences (a different store) +*before* taking the notifications lock, and `checkAndUnlock` applies its unlocks +to the in-flight store object rather than re-entering through +`unlockAchievement`. + +**The same per-process limit in §5.3 applies here.** Serialization and atomic +replacement remove the lost update within one process; they do not reconcile +state across serverless instances. + --- ## 6. On-chain Integration @@ -331,6 +365,10 @@ Rollback: unset the var or set `0` and restart. No ledger migration to revert; o - Use writable server storage abstraction (`server-data-paths`) for platform-safe behavior. - Any store mutated by more than one writer needs serialized read-modify-write, atomic file replacement, and a version/compare-and-swap field. See §5.3. +- Mutate the JSON sidecar stores only through `lib/json-store.ts` + (`updateStore` / `updateStoreWithReader` / `updateJsonFile`). Never pair a raw + `readFile` with a `writeFile`; that is the lost-update race §5.4 exists to + prevent. Locks are not reentrant — see §5.4. - Watch for the `signals.version_conflict` log event to detect clients that are writing against stale reads. - On contract redeploys, update: diff --git a/lib/__tests__/json-store-concurrency.test.ts b/lib/__tests__/json-store-concurrency.test.ts new file mode 100644 index 0000000..b637d1c --- /dev/null +++ b/lib/__tests__/json-store-concurrency.test.ts @@ -0,0 +1,278 @@ +import { createHash } from "node:crypto" +import { mkdtemp, readFile, readdir, writeFile } from "node:fs/promises" +import { tmpdir } from "node:os" +import path from "node:path" +import { Keypair } from "@stellar/stellar-sdk" +import { beforeEach, describe, expect, it } from "vitest" +import { readJsonFile, updateJsonFile, writeJsonFileAtomic } from "@/lib/json-store" +import { serverDataJsonPath } from "@/lib/server-data-paths" +import { followUser, getFollowers, getFollowing, unfollowUser } from "@/lib/follow-store" +import { saveProfile, getProfile } from "@/lib/profile-store" +import { createNotification, createNotificationBatch, getNotifications } from "@/lib/notification-store" +import { checkAndUnlock, getWalletData, unlockAchievement } from "@/lib/achievement-store" +import { saveWorldForCollection, getWorldForCollection, markNarrativeRead, getReaderProgress } from "@/lib/narrative-world-store" +import { recordCreatorProfileView, getCreatorProfileViewAnalytics } from "@/lib/market-store" + +const CONCURRENCY = 50 + +let dataDir = "" + +beforeEach(async () => { + dataDir = await mkdtemp(path.join(tmpdir(), "phase-store-test-")) + process.env.PHASE_SERVER_DATA_DIR = dataDir +}) + +/** + * Deterministic, genuinely valid ed25519 public keys. The market-store profile + * view schema runs `StrKey.isValidEd25519PublicKey`, so synthetic `G...` strings + * are rejected. + */ +const wallet = (n: number) => + Keypair.fromRawEd25519Seed(createHash("sha256").update(`phase-test-${n}`).digest()) + .publicKey() + +describe("follow-store: no lost updates", () => { + it("persists all 50 concurrent distinct follows", async () => { + const target = wallet(500) + await Promise.all( + Array.from({ length: CONCURRENCY }, (_, i) => followUser(wallet(600 + i), target)), + ) + expect((await getFollowers(target)).length).toBe(CONCURRENCY) + }) + + it("records both directions of a follow", async () => { + const a = wallet(700) + const b = wallet(701) + await followUser(a, b) + expect(await getFollowing(a)).toEqual([b]) + expect(await getFollowers(b)).toEqual([a]) + }) + + it("unfollow removes both directions", async () => { + const a = wallet(702) + const b = wallet(703) + await followUser(a, b) + await unfollowUser(a, b) + expect(await getFollowing(a)).toEqual([]) + expect(await getFollowers(b)).toEqual([]) + }) + + it("does not drop a follow when block and follow interleave", async () => { + // Same file, different keys: the unserialized version lost one of the two. + const a = wallet(704) + const b = wallet(705) + const c = wallet(706) + await Promise.all([followUser(a, b), followUser(c, b), followUser(a, c)]) + expect((await getFollowers(b)).sort()).toEqual([a, c].sort()) + expect((await getFollowing(a)).sort()).toEqual([b, c].sort()) + }) +}) + +describe("profile-store: no lost updates", () => { + it("persists all 50 concurrent profile saves", async () => { + await Promise.all( + Array.from({ length: CONCURRENCY }, (_, i) => + saveProfile(wallet(800 + i), { display_name: `name-${i}` }), + ), + ) + for (let i = 0; i < CONCURRENCY; i++) { + expect((await getProfile(wallet(800 + i)))?.display_name).toBe(`name-${i}`) + } + }) +}) + +describe("achievement-store: checkAndUnlock persists what it reports", () => { + it("persists an unlock in a single call (regression: stale snapshot clobber)", async () => { + const w = wallet(1) + const unlocked = await checkAndUnlock(w, { mints: 1 }) + expect(unlocked).toEqual(["first_mint"]) + + const data = await getWalletData(w) + expect(data.unlocked.map((a) => a.id)).toEqual(["first_mint"]) + expect(data.mint_count).toBe(1) + }) + + it("unlockAchievement is idempotent under concurrent calls", async () => { + const w = wallet(2) + const results = await Promise.all( + Array.from({ length: 10 }, () => unlockAchievement(w, "first_collection")), + ) + expect(results.filter(Boolean).length).toBe(1) + expect((await getWalletData(w)).unlocked.length).toBe(1) + }) + + it("accumulates counters across 50 concurrent checkAndUnlock calls", async () => { + const w = wallet(3) + await Promise.all( + Array.from({ length: CONCURRENCY }, () => checkAndUnlock(w, { upvote_delta: 1 })), + ) + const data = await getWalletData(w) + expect(data.total_upvotes).toBe(CONCURRENCY) + expect(data.unlocked.map((a) => a.id)).toEqual(["community_voice"]) + }) + + it("daily_claim streak advances across sequential calls", async () => { + const w = wallet(4) + await checkAndUnlock(w, { daily_claim: true }) + await checkAndUnlock(w, { daily_claim: true }) + expect((await getWalletData(w)).daily_streak).toBe(2) + }) +}) + +describe("notification-store: no lost updates", () => { + it("caps concurrent notifications per wallet without losing the newest", async () => { + const w = wallet(800) + await Promise.all( + Array.from({ length: 60 }, (_, i) => createNotification(w, "signal_reply", { seq: i })), + ) + const list = await getNotifications(w, 100) + expect(list.length).toBe(50) + }) + + it("persists every wallet in a concurrent batch", async () => { + const wallets = Array.from({ length: CONCURRENCY }, (_, i) => wallet(900 + i)) + const res = await createNotificationBatch(wallets, "new_follower", { ok: true }) + expect(res.succeeded).toBe(CONCURRENCY) + expect(res.failed).toBe(0) + for (const w of wallets) { + expect((await getNotifications(w, 10)).length).toBe(1) + } + }) + + it("does not lose notifications when a batch and singles interleave", async () => { + const w = wallet(950) + await Promise.all([ + createNotificationBatch([w], "signal_reply", { batch: true }), + createNotification(w, "signal_upvote", { single: true }), + ]) + const sources = (await getNotifications(w, 10)).map((n) => n.data) + expect(sources).toContainEqual({ batch: true }) + expect(sources).toContainEqual({ single: true }) + }) +}) + +describe("narrative-world-store: version and vector clock are race-free", () => { + it("serializes concurrent world saves into distinct versions", async () => { + const collectionId = 42 + await Promise.all( + Array.from({ length: CONCURRENCY }, (_, i) => + saveWorldForCollection(collectionId, { + world_name: `world-${i}`, + world_prompt: "p", + }), + ), + ) + const saved = await getWorldForCollection(collectionId) + // Every writer must have observed a distinct prior version, so the counter + // has to reach the number of writers rather than collapsing to 1. + expect(saved?.version).toBe(CONCURRENCY) + }) + + it("accumulates concurrent reader progress", async () => { + const w = wallet(1000) + const collectionId = 7 + await Promise.all( + Array.from({ length: 20 }, (_, i) => markNarrativeRead(w, collectionId, i)), + ) + const read = await getReaderProgress(w, collectionId) + expect(read.length).toBe(20) + }) +}) + +describe("market-store: profile view counters", () => { + it("counts all 50 concurrent profile views", async () => { + const creator = wallet(1100) + await Promise.all( + Array.from({ length: CONCURRENCY }, (_, i) => + recordCreatorProfileView( + { creator_wallet: creator, viewer_wallet: wallet(1200 + i), source: "profile" }, + { force: true }, + ), + ), + ) + const analytics = await getCreatorProfileViewAnalytics(creator) + expect(analytics?.total_views).toBe(CONCURRENCY) + expect(analytics?.unique_viewers).toBe(CONCURRENCY) + }) +}) + +describe("json-store primitives", () => { + it("atomic write leaves no temp files behind", async () => { + const file = path.join(dataDir, "atomic.json") + await writeJsonFileAtomic(file, { ok: true }) + expect(JSON.parse(await readFile(file, "utf8"))).toEqual({ ok: true }) + const entries = await readdir(dataDir) + expect(entries.filter((e) => e.endsWith(".tmp"))).toEqual([]) + }) + + it("atomic write overwrites an existing file completely", async () => { + const file = path.join(dataDir, "atomic.json") + await writeJsonFileAtomic(file, { big: "x".repeat(50_000) }) + await writeJsonFileAtomic(file, { small: true }) + expect(JSON.parse(await readFile(file, "utf8"))).toEqual({ small: true }) + }) + + it("returns the fallback for a missing file", async () => { + const missing = path.join(dataDir, "nope.json") + expect(await readJsonFile(missing, { fallback: true })).toEqual({ fallback: true }) + }) + + it("throws on a corrupt store instead of degrading to empty", async () => { + const file = path.join(dataDir, "corrupt.json") + await writeFile(file, "{ not json", "utf8") + await expect(readJsonFile(file, {})).rejects.toThrow(/Corrupt JSON store/) + }) + + it("does not poison the lock chain when a mutation throws", async () => { + const file = path.join(dataDir, "chain.json") + await expect( + updateJsonFile<{ n: number }, void>(file, { + mutate: () => { + throw new Error("boom") + }, + }), + ).rejects.toThrow(/boom/) + + await updateJsonFile<{ n: number }, void>(file, { + mutate: (s) => { + s.n = 1 + }, + }) + expect(JSON.parse(await readFile(file, "utf8"))).toEqual({ n: 1 }) + }) + + it("serializes overlapping mutations on the same path", async () => { + const file = path.join(dataDir, "serial.json") + const order: string[] = [] + await Promise.all( + Array.from({ length: 5 }, (_, i) => + updateJsonFile<{ seen: number[] }, void>(file, { + mutate: async (s) => { + order.push(`enter-${i}`) + await new Promise((r) => setTimeout(r, 5)) + s.seen = [...(s.seen ?? []), i] + order.push(`exit-${i}`) + }, + }), + ), + ) + for (let i = 0; i < order.length; i += 2) { + expect(order[i].replace("enter-", "")).toBe(order[i + 1].replace("exit-", "")) + } + expect(JSON.parse(await readFile(file, "utf8")).seen.length).toBe(5) + }) + + it("applies the read adapter before mutating", async () => { + const file = path.join(dataDir, "adapter.json") + await writeJsonFileAtomic(file, { legacy: 7 }) + const seen = await updateJsonFile<{ value: number }, number>(file, { + read: (raw) => ({ value: (raw as { legacy?: number }).legacy ?? 0 }), + mutate: (s) => { + s.value += 1 + return s.value + }, + }) + expect(seen).toBe(8) + expect(JSON.parse(await readFile(file, "utf8"))).toEqual({ value: 8 }) + }) +}) diff --git a/lib/achievement-store.ts b/lib/achievement-store.ts index e7c0966..d857f8e 100644 --- a/lib/achievement-store.ts +++ b/lib/achievement-store.ts @@ -1,5 +1,4 @@ -import { mkdir, readFile, writeFile } from "node:fs/promises"; -import path from "node:path"; +import { readJsonFile, updateStore, updateStoreWithReader } from "@/lib/json-store"; import { serverDataJsonPath } from "@/lib/server-data-paths"; import { createNotification } from "@/lib/notification-store"; @@ -79,27 +78,33 @@ function mergeEntries( }; } -async function readStore(): Promise { - let raw: AchievementStore; - try { - raw = JSON.parse( - await readFile(serverDataJsonPath("achievements"), "utf8"), - ) as AchievementStore; - } catch { - return {}; - } +/** Normalizes wallet keys and merges duplicate-cased rows, as the old reader did. */ +function readAchievementStore(raw: unknown): AchievementStore { + if (!raw || typeof raw !== "object") return {}; const store: AchievementStore = {}; - for (const [wallet, entry] of Object.entries(raw)) { + for (const [wallet, entry] of Object.entries(raw as AchievementStore)) { const key = walletKey(wallet); store[key] = store[key] ? mergeEntries(store[key]!, entry) : entry; } return store; } -async function writeStore(data: AchievementStore): Promise { - const fp = serverDataJsonPath("achievements"); - await mkdir(path.dirname(fp), { recursive: true }); - await writeFile(fp, JSON.stringify(data, null, 2), "utf8"); +async function readStore(): Promise { + return readJsonFile( + serverDataJsonPath("achievements"), + undefined, + ).then(readAchievementStore); +} + +/** Locked read-modify-write that normalizes wallet keys before mutating. */ +function updateAchievementStore( + mutate: (store: AchievementStore) => R | Promise, +): Promise { + return updateStoreWithReader( + "achievements", + readAchievementStore, + mutate, + ); } function ensureEntry( @@ -128,20 +133,23 @@ export async function unlockAchievement( id: AchievementId, evidence?: string, ): Promise { - const store = await readStore(); - const entry = ensureEntry(store, wallet); - if (entry.unlocked.some((a) => a.id === id)) return false; // idempotent - entry.unlocked.push({ id, unlocked_at: Date.now(), tx_evidence: evidence }); - store[walletKey(wallet)] = entry; - await writeStore(store); - // Notify (fire-and-forget) - void createNotification(wallet, "achievement_unlocked", { - achievement_id: id, - achievement_name: ACHIEVEMENT_NAMES[id] ?? id, - }).catch(() => { - /* silent */ + const didUnlock = await updateAchievementStore((store) => { + const entry = ensureEntry(store, wallet); + if (entry.unlocked.some((a) => a.id === id)) return false; // idempotent + entry.unlocked.push({ id, unlocked_at: Date.now(), tx_evidence: evidence }); + store[walletKey(wallet)] = entry; + return true; }); - return true; + if (didUnlock) { + // Notify (fire-and-forget) after the store lock is released. + void createNotification(wallet, "achievement_unlocked", { + achievement_id: id, + achievement_name: ACHIEVEMENT_NAMES[id] ?? id, + }).catch(() => { + /* silent */ + }); + } + return didUnlock; } /** Checks counters and unlocks newly earned achievements. Returns newly unlocked IDs. */ @@ -159,75 +167,90 @@ export async function checkAndUnlock( phaselq_earned?: number; }, ): Promise { - const store = await readStore(); - const entry = ensureEntry(store, wallet); - const unlocked = new Set(entry.unlocked.map((a) => a.id)); - const newUnlocks: AchievementId[] = []; - - async function tryUnlock(id: AchievementId, evidence?: string) { - if (unlocked.has(id)) return; - const didUnlock = await unlockAchievement(wallet, id, evidence); - if (didUnlock) { - newUnlocks.push(id); + // The whole counter update and every unlock it triggers happen in ONE locked + // mutation. Previously this read the store, let tryUnlock call + // unlockAchievement (its own read-modify-write), then wrote this stale + // snapshot back over the top — so checkAndUnlock reported unlocks to its + // caller and sent "achievement unlocked" notifications for achievements that + // were never persisted. Re-entering the store lock from inside a locked + // mutation would deadlock, so unlocks are applied to the in-flight store. + const newUnlocks = await updateAchievementStore((store) => { + const entry = ensureEntry(store, wallet); + const unlocked = new Set(entry.unlocked.map((a) => a.id)); + const newlyUnlocked: AchievementId[] = []; + + function tryUnlock(id: AchievementId) { + if (unlocked.has(id)) return; + entry.unlocked.push({ id, unlocked_at: Date.now() }); unlocked.add(id); + newlyUnlocked.push(id); } - } - // Mint counts - if (hints?.mints !== undefined) { - entry.mint_count = (entry.mint_count ?? 0) + hints.mints; - if (entry.mint_count >= 1) await tryUnlock("first_mint"); - if (entry.mint_count >= 5) await tryUnlock("collector_5"); - if (entry.mint_count >= 10) await tryUnlock("collector_10"); - } + // Mint counts + if (hints?.mints !== undefined) { + entry.mint_count = (entry.mint_count ?? 0) + hints.mints; + if (entry.mint_count >= 1) tryUnlock("first_mint"); + if (entry.mint_count >= 5) tryUnlock("collector_5"); + if (entry.mint_count >= 10) tryUnlock("collector_10"); + } - // First collection - if (hints?.has_collection) await tryUnlock("first_collection"); + // First collection + if (hints?.has_collection) tryUnlock("first_collection"); - // World builder - if (hints?.has_world) await tryUnlock("world_builder"); + // World builder + if (hints?.has_world) tryUnlock("world_builder"); - // Signal pioneer - if (hints?.signal_posted) await tryUnlock("signal_pioneer"); + // Signal pioneer + if (hints?.signal_posted) tryUnlock("signal_pioneer"); - // Upvotes - if (hints?.upvote_delta !== undefined) { - entry.total_upvotes = (entry.total_upvotes ?? 0) + hints.upvote_delta; - if (entry.total_upvotes >= 25) await tryUnlock("community_voice"); - } + // Upvotes + if (hints?.upvote_delta !== undefined) { + entry.total_upvotes = (entry.total_upvotes ?? 0) + hints.upvote_delta; + if (entry.total_upvotes >= 25) tryUnlock("community_voice"); + } - // Followers - if (hints?.follower_delta !== undefined) { - entry.follower_count = (entry.follower_count ?? 0) + hints.follower_delta; - if (entry.follower_count >= 10) await tryUnlock("connector_10"); - } + // Followers + if (hints?.follower_delta !== undefined) { + entry.follower_count = (entry.follower_count ?? 0) + hints.follower_delta; + if (entry.follower_count >= 10) tryUnlock("connector_10"); + } - // Narrator - if (hints?.narrator_delta !== undefined) { - entry.narrator_count = (entry.narrator_count ?? 0) + hints.narrator_delta; - if (entry.narrator_count >= 10) await tryUnlock("narrator_10"); - } + // Narrator + if (hints?.narrator_delta !== undefined) { + entry.narrator_count = (entry.narrator_count ?? 0) + hints.narrator_delta; + if (entry.narrator_count >= 10) tryUnlock("narrator_10"); + } - // Daily streak - if (hints?.daily_claim) { - const now = Date.now(); - const last = entry.last_daily ?? 0; - const dayMs = 86_400_000; - const withinWindow = last > 0 && now - last < dayMs * 2; - entry.daily_streak = withinWindow ? (entry.daily_streak ?? 0) + 1 : 1; - entry.last_daily = now; - if (entry.daily_streak >= 7) await tryUnlock("daily_streak_7"); - if (entry.daily_streak >= 30) await tryUnlock("daily_streak_30"); - } + // Daily streak + if (hints?.daily_claim) { + const now = Date.now(); + const last = entry.last_daily ?? 0; + const dayMs = 86_400_000; + const withinWindow = last > 0 && now - last < dayMs * 2; + entry.daily_streak = withinWindow ? (entry.daily_streak ?? 0) + 1 : 1; + entry.last_daily = now; + if (entry.daily_streak >= 7) tryUnlock("daily_streak_7"); + if (entry.daily_streak >= 30) tryUnlock("daily_streak_30"); + } + + // PHASELQ (placeholder — would need tracking from faucet totals) + // For now just check if they've earned any + if (hints?.phaselq_earned !== undefined && hints.phaselq_earned >= 100) { + tryUnlock("phaselq_100"); + } + + store[walletKey(wallet)] = entry; + return newlyUnlocked; + }); - // PHASELQ (placeholder — would need tracking from faucet totals) - // For now just check if they've earned any - if (hints?.phaselq_earned !== undefined && hints.phaselq_earned >= 100) { - await tryUnlock("phaselq_100"); + // Notify after the store lock is released, matching unlockAchievement. + for (const id of newUnlocks) { + void createNotification(wallet, "achievement_unlocked", { + achievement_id: id, + achievement_name: ACHIEVEMENT_NAMES[id] ?? id, + }).catch(() => { /* silent */ }); } - store[walletKey(wallet)] = entry; - await writeStore(store); return newUnlocks; } diff --git a/lib/follow-store.ts b/lib/follow-store.ts index dbebb23..2b46232 100644 --- a/lib/follow-store.ts +++ b/lib/follow-store.ts @@ -1,7 +1,6 @@ -import { mkdir, readFile, writeFile } from "node:fs/promises"; -import path from "node:path"; import { z } from "zod"; import { isFeatureEnabled, flagRollbackNote } from "@/lib/feature-flags"; +import { readJsonFile, updateStore } from "@/lib/json-store"; import { serverDataJsonPath } from "@/lib/server-data-paths"; import { HORIZON_URL } from "@/lib/phase-protocol"; @@ -220,19 +219,7 @@ export function validateSep50MetadataBeforePin( } async function readStore(): Promise { - try { - return JSON.parse( - await readFile(serverDataJsonPath("profileFollows"), "utf8"), - ) as FollowStore; - } catch { - return {}; - } -} - -async function writeStore(data: FollowStore): Promise { - const filePath = serverDataJsonPath("profileFollows"); - await mkdir(path.dirname(filePath), { recursive: true }); - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8"); + return readJsonFile(serverDataJsonPath("profileFollows"), {}); } function ensureEntry(store: FollowStore, wallet: string): FollowEntry { @@ -245,24 +232,24 @@ export async function followUser( toWallet: string, ): Promise { if (fromWallet === toWallet) return; - const store = await readStore(); - const from = ensureEntry(store, fromWallet); - const to = ensureEntry(store, toWallet); - if (!from.following.includes(toWallet)) from.following.push(toWallet); - if (!to.followers.includes(fromWallet)) to.followers.push(fromWallet); - await writeStore(store); + await updateStore("profileFollows", (store) => { + const from = ensureEntry(store, fromWallet); + const to = ensureEntry(store, toWallet); + if (!from.following.includes(toWallet)) from.following.push(toWallet); + if (!to.followers.includes(fromWallet)) to.followers.push(fromWallet); + }); } export async function unfollowUser( fromWallet: string, toWallet: string, ): Promise { - const store = await readStore(); - const from = ensureEntry(store, fromWallet); - const to = ensureEntry(store, toWallet); - from.following = from.following.filter((w) => w !== toWallet); - to.followers = to.followers.filter((w) => w !== fromWallet); - await writeStore(store); + await updateStore("profileFollows", (store) => { + const from = ensureEntry(store, fromWallet); + const to = ensureEntry(store, toWallet); + from.following = from.following.filter((w) => w !== toWallet); + to.followers = to.followers.filter((w) => w !== fromWallet); + }); } export async function getFollowers(wallet: string): Promise { diff --git a/lib/json-store.ts b/lib/json-store.ts new file mode 100644 index 0000000..34d7f19 --- /dev/null +++ b/lib/json-store.ts @@ -0,0 +1,167 @@ +import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; +import path from "node:path"; +import { serverDataJsonPath, type ServerDataFile } from "@/lib/server-data-paths"; + +/** + * Serialized, crash-safe access to the JSON sidecar stores. + * + * Each store used to hand-roll its own `readFile` -> `JSON.parse` -> mutate -> + * `writeFile` pair with no coordination, so two concurrent read-modify-write + * cycles on the same file raced and the second `writeFile` silently discarded + * the first. `lib/signal-store.ts` had this fixed inline by #361; the remaining + * JSON stores share the same shape and route their mutations through here. + * + * Three guarantees: + * + * 1. `withFileLock` chains work per absolute file path, so a read-modify-write + * cycle runs to completion before the next one starts. `updateJsonFile` puts + * the read AND the write inside that chain, which is what removes the lost + * update. + * 2. `writeJsonFileAtomic` writes a sibling temp file then `rename`s it over the + * target, so a reader never observes a half-written file and a crash + * mid-write cannot truncate the store to invalid JSON. + * 3. A JSON parse failure throws instead of degrading to `{}`. Returning `{}` + * on a parse error meant the next write silently overwrote every existing + * record. A missing file still resolves to the caller's fallback. + * + * SCOPE LIMIT — the lock is per-process. It removes the single-instance race, + * which is the dominant failure mode, but it does NOT make the sidecars safe + * across processes or serverless instances: on Vercel each instance resolves + * its own `os.tmpdir()` copy of the data root, so two instances hold divergent + * state and whichever writes last wins. Cross-instance consistency needs a + * shared datastore or an advisory lock on a shared volume. See + * `serverDataRoot()` in lib/server-data-paths.ts and docs/TECHNICAL.md 5.4. + * + * LOCKS ARE NOT REENTRANT. A `mutate` callback must never call back into a + * locking store function for the same file, or it will deadlock waiting on its + * own chain. Apply nested changes to the in-flight `store` object directly and + * fire notifications after the callback returns. + */ + +const fileQueues = new Map>(); + +/** + * Runs `task` with exclusive access to `filePath` within this process. Pairing + * this with the raw read/write helpers is almost never what you want — use + * `updateJsonFile` / `updateStore` so the read and the write share one critical + * section. + */ +export function withFileLock( + filePath: string, + task: () => Promise, +): Promise { + const previous = fileQueues.get(filePath) ?? Promise.resolve(); + const next = previous.then(task, task); + fileQueues.set( + filePath, + next.then( + () => undefined, + () => undefined, + ), + ); + return next; +} + +function isErrnoCode(error: unknown, code: string): boolean { + return ( + typeof error === "object" && + error !== null && + (error as NodeJS.ErrnoException).code === code + ); +} + +/** + * Unlocked read. A missing file yields `fallback`; a corrupt file throws rather + * than silently reporting empty. Safe for read-only consumers. + */ +export async function readJsonFile( + filePath: string, + fallback: T, +): Promise { + let raw: string; + try { + raw = await readFile(filePath, "utf8"); + } catch (error) { + if (isErrnoCode(error, "ENOENT")) return fallback; + throw error; + } + try { + return JSON.parse(raw) as T; + } catch (error) { + throw new Error(`Corrupt JSON store at ${filePath}: ${String(error)}`); + } +} + +/** + * Unlocked atomic write. Prefer `updateJsonFile` / `updateStore` — this only + * guarantees the write is not torn, not that it will not clobber a concurrent + * writer. + */ +export async function writeJsonFileAtomic( + filePath: string, + data: unknown, +): Promise { + await mkdir(path.dirname(filePath), { recursive: true }); + const tmpPath = `${filePath}.${process.pid}.${Date.now().toString(36)}.tmp`; + try { + await writeFile(tmpPath, JSON.stringify(data, null, 2), "utf8"); + await rename(tmpPath, filePath); + } catch (error) { + await rm(tmpPath, { force: true }).catch(() => undefined); + throw error; + } +} + +/** + * Locked read-modify-write: reads, applies `mutate`, and atomically writes while + * holding the file's lock, then returns whatever `mutate` returned. This is the + * only correct way to mutate a store. + * + * `read` adapts the parsed JSON into the store's shape. Use it for files with a + * legacy or defensive on-disk format that needs normalizing before mutation. + */ +export async function updateJsonFile( + filePath: string, + opts: { + read?: (raw: unknown) => T; + mutate: (store: T) => R | Promise; + }, +): Promise { + return withFileLock(filePath, async () => { + const raw = await readJsonFile(filePath, undefined); + // A missing file starts from an empty store, matching the + // `catch { return {} }` readers this replaces. + const store = opts.read ? opts.read(raw) : ((raw ?? {}) as T); + const result = await opts.mutate(store); + await writeJsonFileAtomic(filePath, store); + return result; + }); +} + +/** Locked read-modify-write for a registered sidecar file. */ +export function updateStore( + key: ServerDataFile, + mutate: (store: T) => R | Promise, +): Promise { + return updateJsonFile(serverDataJsonPath(key), { mutate }); +} + +/** + * Locked read-modify-write for a registered sidecar file whose on-disk shape + * needs normalizing before mutation. + */ +export function updateStoreWithReader( + key: ServerDataFile, + read: (raw: unknown) => T, + mutate: (store: T) => R | Promise, +): Promise { + return updateJsonFile(serverDataJsonPath(key), { read, mutate }); +} + +/** Unlocked read of a registered sidecar file. Safe for read-only consumers. */ +export function readStore( + key: ServerDataFile, + fallback: () => T, +): Promise { + return readJsonFile(serverDataJsonPath(key), fallback()); +} diff --git a/lib/market-store.ts b/lib/market-store.ts index d9334da..12d89f4 100644 --- a/lib/market-store.ts +++ b/lib/market-store.ts @@ -1,10 +1,9 @@ // @ts-nocheck -import { mkdir, readFile, writeFile } from "node:fs/promises"; -import path from "node:path"; import { createHash, randomUUID } from "node:crypto"; import { StrKey } from "@stellar/stellar-sdk"; import { z } from "zod"; import { isFeatureEnabled, flagRollbackNote } from "@/lib/feature-flags"; +import { readJsonFile, updateJsonFile } from "@/lib/json-store"; import { serverDataJsonPath } from "@/lib/server-data-paths"; import { getDb } from "@/lib/sqlite-db"; import { incSecurityCounter } from "@/lib/security-counters"; @@ -100,19 +99,7 @@ export function phase100RollbackNote(): string { } async function readJson(filePath: string): Promise { - try { - return JSON.parse(await readFile(filePath, "utf8")) as T; - } catch { - return {} as T; - } -} - -async function writeJson( - filePath: string, - data: T, -): Promise { - await mkdir(path.dirname(filePath), { recursive: true }); - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8"); + return readJsonFile(filePath, {} as T); } function viewerAnalyticsKey(viewerWallet?: string): string | null { @@ -147,39 +134,45 @@ export async function recordCreatorProfileView( const event = parsed.data; const now = opts.now ?? Date.now(); - const store = await readJson( - serverDataJsonPath("marketProfileViews"), - ); - const current = store[event.creator_wallet] ?? { - creator_wallet: event.creator_wallet, - total_views: 0, - unique_viewers: 0, - last_viewed_at: 0, - sources: {}, - viewer_hashes: [], - }; - const viewerKey = viewerAnalyticsKey(event.viewer_wallet); - const viewerHashes = - viewerKey && !current.viewer_hashes.includes(viewerKey) - ? [...current.viewer_hashes, viewerKey] - : current.viewer_hashes; - - const next: CreatorProfileViewAnalytics = { - ...current, - total_views: current.total_views + 1, - unique_viewers: viewerHashes.length, - last_viewed_at: now, - sources: { - ...current.sources, - [event.source]: (current.sources[event.source] ?? 0) + 1, + // The counters are derived from the snapshot read here, so this has to be one + // locked mutation: an unguarded read-then-write dropped concurrent views. + return updateJsonFile( + serverDataJsonPath("marketProfileViews"), + { + mutate: (store) => { + const current = store[event.creator_wallet] ?? { + creator_wallet: event.creator_wallet, + total_views: 0, + unique_viewers: 0, + last_viewed_at: 0, + sources: {}, + viewer_hashes: [], + }; + + const viewerKey = viewerAnalyticsKey(event.viewer_wallet); + const viewerHashes = + viewerKey && !current.viewer_hashes.includes(viewerKey) + ? [...current.viewer_hashes, viewerKey] + : current.viewer_hashes; + + const next: CreatorProfileViewAnalytics = { + ...current, + total_views: current.total_views + 1, + unique_viewers: viewerHashes.length, + last_viewed_at: now, + sources: { + ...current.sources, + [event.source]: (current.sources[event.source] ?? 0) + 1, + }, + viewer_hashes: viewerHashes, + }; + + store[event.creator_wallet] = next; + return next; + }, }, - viewer_hashes: viewerHashes, - }; - - store[event.creator_wallet] = next; - await writeJson(serverDataJsonPath("marketProfileViews"), store); - return next; + ); } export async function getCreatorProfileViewAnalytics( @@ -795,13 +788,15 @@ export async function blockWallet( reason?: string, ): Promise { if (!isPhase85Enabled()) throw new Error("phase-85 disabled"); - const store = await readJson(serverDataJsonPath("blockList")); - const list = store[blocker] ?? { blocked: [], muted: [] }; - if (!list.blocked.some((b) => b.wallet === target)) { - list.blocked.push({ wallet: target, blocked_at: Date.now(), reason }); - } - store[blocker] = list; - await writeJson(serverDataJsonPath("blockList"), store); + await updateJsonFile(serverDataJsonPath("blockList"), { + mutate: (store) => { + const list = store[blocker] ?? { blocked: [], muted: [] }; + if (!list.blocked.some((b) => b.wallet === target)) { + list.blocked.push({ wallet: target, blocked_at: Date.now(), reason }); + } + store[blocker] = list; + }, + }); } export async function unblockWallet( @@ -809,11 +804,13 @@ export async function unblockWallet( target: string, ): Promise { if (!isPhase85Enabled()) throw new Error("phase-85 disabled"); - const store = await readJson(serverDataJsonPath("blockList")); - const list = store[blocker] ?? { blocked: [], muted: [] }; - list.blocked = list.blocked.filter((b) => b.wallet !== target); - store[blocker] = list; - await writeJson(serverDataJsonPath("blockList"), store); + await updateJsonFile(serverDataJsonPath("blockList"), { + mutate: (store) => { + const list = store[blocker] ?? { blocked: [], muted: [] }; + list.blocked = list.blocked.filter((b) => b.wallet !== target); + store[blocker] = list; + }, + }); } export async function muteWallet( @@ -822,16 +819,18 @@ export async function muteWallet( durationMs?: number, ): Promise { if (!isPhase85Enabled()) throw new Error("phase-85 disabled"); - const store = await readJson(serverDataJsonPath("blockList")); - const list = store[muter] ?? { blocked: [], muted: [] }; - list.muted = list.muted.filter((m) => m.wallet !== target); - list.muted.push({ - wallet: target, - muted_at: Date.now(), - expires_at: durationMs ? Date.now() + durationMs : undefined, + await updateJsonFile(serverDataJsonPath("blockList"), { + mutate: (store) => { + const list = store[muter] ?? { blocked: [], muted: [] }; + list.muted = list.muted.filter((m) => m.wallet !== target); + list.muted.push({ + wallet: target, + muted_at: Date.now(), + expires_at: durationMs ? Date.now() + durationMs : undefined, + }); + store[muter] = list; + }, }); - store[muter] = list; - await writeJson(serverDataJsonPath("blockList"), store); } export async function unmuteWallet( @@ -839,11 +838,13 @@ export async function unmuteWallet( target: string, ): Promise { if (!isPhase85Enabled()) throw new Error("phase-85 disabled"); - const store = await readJson(serverDataJsonPath("blockList")); - const list = store[muter] ?? { blocked: [], muted: [] }; - list.muted = list.muted.filter((m) => m.wallet !== target); - store[muter] = list; - await writeJson(serverDataJsonPath("blockList"), store); + await updateJsonFile(serverDataJsonPath("blockList"), { + mutate: (store) => { + const list = store[muter] ?? { blocked: [], muted: [] }; + list.muted = list.muted.filter((m) => m.wallet !== target); + store[muter] = list; + }, + }); } export async function isWalletBlocked( diff --git a/lib/narrative-world-store.ts b/lib/narrative-world-store.ts index 669b78b..c9beabc 100644 --- a/lib/narrative-world-store.ts +++ b/lib/narrative-world-store.ts @@ -1,6 +1,5 @@ // @ts-nocheck -import { mkdir, readFile, writeFile } from "node:fs/promises" -import path from "node:path" +import { readJsonFile, updateJsonFile } from "@/lib/json-store" import { serverDataJsonPath } from "@/lib/server-data-paths" import { incrementVectorClock, type VectorClock } from "@/lib/world-conflict" @@ -27,17 +26,20 @@ type WorldCollectionsStore = Record type WorldNarrativesStore = Record async function readJsonStore(filePath: string): Promise { - try { - const raw = await readFile(filePath, "utf8") - return JSON.parse(raw) as T - } catch { - return {} as T - } + return readJsonFile(filePath, {} as T) } -async function writeJsonStore(filePath: string, data: T): Promise { - await mkdir(path.dirname(filePath), { recursive: true }) - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8") +/** + * Locked read-modify-write. Required for every mutation here: these stores + * compute counters, versions and vector clocks from the snapshot they read, so + * an unguarded read-then-write silently loses a concurrent writer's contribution + * and can hand two writers the same vector clock. + */ +async function updateJsonStore( + filePath: string, + mutate: (store: T) => R | Promise, +): Promise { + return updateJsonFile(filePath, { mutate }) } export async function getWorldForCollection(collectionId: number): Promise { @@ -53,22 +55,24 @@ export async function saveWorldForCollection( creator_wallet?: string }, ): Promise { - const filePath = serverDataJsonPath("worldCollections") - const store = await readJsonStore(filePath) - const existing = store[String(collectionId)] - const saved: WorldCollectionData = { - ...existing, - world_name: data.world_name, - world_prompt: data.world_prompt, - ...(data.narrator_tone !== undefined ? { narrator_tone: data.narrator_tone } : {}), - ...(data.creator_wallet !== undefined ? { creator_wallet: data.creator_wallet } : {}), - created_at: existing?.created_at ?? Date.now(), - version: (existing?.version ?? 0) + 1, - vector_clock: incrementVectorClock(existing?.vector_clock, data.creator_wallet ?? "anonymous"), - } - store[String(collectionId)] = saved - await writeJsonStore(filePath, store) - return saved + return updateJsonStore( + serverDataJsonPath("worldCollections"), + (store) => { + const existing = store[String(collectionId)] + const saved: WorldCollectionData = { + ...existing, + world_name: data.world_name, + world_prompt: data.world_prompt, + ...(data.narrator_tone !== undefined ? { narrator_tone: data.narrator_tone } : {}), + ...(data.creator_wallet !== undefined ? { creator_wallet: data.creator_wallet } : {}), + created_at: existing?.created_at ?? Date.now(), + version: (existing?.version ?? 0) + 1, + vector_clock: incrementVectorClock(existing?.vector_clock, data.creator_wallet ?? "anonymous"), + } + store[String(collectionId)] = saved + return saved + }, + ) } export async function getAllWorldCollections(): Promise { @@ -86,10 +90,12 @@ export async function saveNarrativeForToken( tokenId: number, data: Omit, ): Promise { - const filePath = serverDataJsonPath("worldNarratives") - const store = await readJsonStore(filePath) - store[String(tokenId)] = { ...data, generated_at: Date.now() } - await writeJsonStore(filePath, store) + await updateJsonStore( + serverDataJsonPath("worldNarratives"), + (store) => { + store[String(tokenId)] = { ...data, generated_at: Date.now() } + }, + ) invalidateLocalizedNarrativeCache(tokenId) } @@ -221,22 +227,24 @@ function readerProgressKey(wallet: string, collectionId: number): string { } export async function getReaderProgress(wallet: string, collectionId: number): Promise { - const store = await readJsonStore(serverDataJsonPath("readerProgress")) + const store = await readJsonStore(serverDataJsonPath("worldReaderProgress")) const key = readerProgressKey(wallet, collectionId) return store[key]?.read_token_ids ?? [] } export async function markNarrativeRead(wallet: string, collectionId: number, tokenId: number): Promise { - const filePath = serverDataJsonPath("readerProgress") - const store = await readJsonStore(filePath) - const key = readerProgressKey(wallet, collectionId) - const existing = store[key] ?? { wallet, collection_id: collectionId, read_token_ids: [], last_read_at: 0 } - if (!existing.read_token_ids.includes(tokenId)) { - existing.read_token_ids.push(tokenId) - } - existing.last_read_at = Date.now() - store[key] = existing - await writeJsonStore(filePath, store) + await updateJsonStore( + serverDataJsonPath("worldReaderProgress"), + (store) => { + const key = readerProgressKey(wallet, collectionId) + const existing = store[key] ?? { wallet, collection_id: collectionId, read_token_ids: [], last_read_at: 0 } + if (!existing.read_token_ids.includes(tokenId)) { + existing.read_token_ids.push(tokenId) + } + existing.last_read_at = Date.now() + store[key] = existing + }, + ) } // ─── phase-109: collaborative world permissions ────────────────────────── @@ -257,13 +265,12 @@ export async function getWorldRoles(collectionId: number): Promise { - const filePath = serverDataJsonPath("worldRoles") - const store = await readJsonStore(filePath) - const key = String(collectionId) - if (!store[key]) { - store[key] = { collection_id: collectionId, owner: ownerWallet, roles: {} } - await writeJsonStore(filePath, store) - } + await updateJsonStore(serverDataJsonPath("worldRoles"), (store) => { + const key = String(collectionId) + if (!store[key]) { + store[key] = { collection_id: collectionId, owner: ownerWallet, roles: {} } + } + }) } export async function setWorldRole( @@ -272,16 +279,18 @@ export async function setWorldRole( targetWallet: string, role: WorldRole, ): Promise> { - const filePath = serverDataJsonPath("worldRoles") - const store = await readJsonStore(filePath) - const key = String(collectionId) - const entry = store[key] - if (!entry || entry.owner !== actingWallet) { - throw new Error("Solo el propietario del mundo puede asignar roles") - } - entry.roles[targetWallet] = role - await writeJsonStore(filePath, store) - return entry.roles + return updateJsonStore>( + serverDataJsonPath("worldRoles"), + (store) => { + const key = String(collectionId) + const entry = store[key] + if (!entry || entry.owner !== actingWallet) { + throw new Error("Solo el propietario del mundo puede asignar roles") + } + entry.roles[targetWallet] = role + return entry.roles + }, + ) } // ─── phase-112: world export to portable markdown/JSON ─────────────────── @@ -476,22 +485,24 @@ export async function getLoreLinksForToken(tokenId: number): Promise<{ outgoing: } export async function addLoreLink(fromTokenId: number, toTokenId: number, note?: string): Promise { - const filePath = serverDataJsonPath("loreLinks") - const store = await readJsonStore(filePath) - const existingIndex = store.findIndex((l) => l.from_token_id === fromTokenId && l.to_token_id === toTokenId) - const link: LoreLink = { - from_token_id: fromTokenId, - to_token_id: toTokenId, - note, - created_at: Date.now(), - } - if (existingIndex >= 0) { - store[existingIndex] = link - } else { - store.push(link) - } - await writeJsonStore(filePath, store) - return link + return updateJsonStore( + serverDataJsonPath("loreLinks"), + (store) => { + const existingIndex = store.findIndex((l) => l.from_token_id === fromTokenId && l.to_token_id === toTokenId) + const link: LoreLink = { + from_token_id: fromTokenId, + to_token_id: toTokenId, + note, + created_at: Date.now(), + } + if (existingIndex >= 0) { + store[existingIndex] = link + } else { + store.push(link) + } + return link + }, + ) } // ─── phase-110: narrative search helpers ────────────────────────────────── diff --git a/lib/notification-store.ts b/lib/notification-store.ts index 25a8f16..1f33702 100644 --- a/lib/notification-store.ts +++ b/lib/notification-store.ts @@ -1,9 +1,8 @@ -import { mkdir, readFile, writeFile } from "node:fs/promises" -import path from "node:path" import { randomUUID } from "node:crypto" import { createHash } from "node:crypto" import { z } from "zod" import { isFeatureEnabled, flagRollbackNote } from "@/lib/feature-flags" +import { readJsonFile, updateStore, updateStoreWithReader } from "@/lib/json-store" import { serverDataJsonPath } from "@/lib/server-data-paths" export type NotificationType = @@ -106,20 +105,12 @@ function gatewayRotationKey(gateway: string, privateTier: string): string { } async function readGatewayAuthRotationStore(): Promise { - try { - const raw = await readFile(serverDataJsonPath("ipfsGatewayAuthRotations"), "utf8") - const parsed = JSON.parse(raw) as GatewayAuthRotationStore - return parsed && typeof parsed === "object" ? parsed : {} - } catch { - return {} - } + return readJsonFile( + serverDataJsonPath("ipfsGatewayAuthRotations"), + {}, + ) } -async function writeGatewayAuthRotationStore(data: GatewayAuthRotationStore): Promise { - const filePath = serverDataJsonPath("ipfsGatewayAuthRotations") - await mkdir(path.dirname(filePath), { recursive: true }) - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8") -} export async function rotateIpfsGatewayAuth(input: unknown, opts: { force?: boolean; now?: number } = {}): Promise { if (!opts.force && !isPhase128Enabled()) { @@ -135,45 +126,40 @@ export async function rotateIpfsGatewayAuth(input: unknown, opts: { force?: bool const now = opts.now ?? Date.now() const data = parsed.data - const store = await readGatewayAuthRotationStore() const key = gatewayRotationKey(data.gateway, data.private_tier) - const current = store[key] - const activeTokenHash = hashGatewayToken(data.next_token) - - const rotation: GatewayAuthRotation = { - gateway: data.gateway, - private_tier: data.private_tier, - active_token_hash: activeTokenHash, - previous_token_hash: current?.active_token_hash && current.active_token_hash !== activeTokenHash - ? current.active_token_hash - : current?.previous_token_hash ?? null, - previous_expires_at: current?.active_token_hash && current.active_token_hash !== activeTokenHash - ? now + data.overlap_ms - : current?.previous_expires_at ?? null, - rotated_by: data.rotated_by, - rotated_at: now, - } + return updateStore( + "ipfsGatewayAuthRotations", + (store) => { + const current = store[key] + const activeTokenHash = hashGatewayToken(data.next_token) + + const rotation: GatewayAuthRotation = { + gateway: data.gateway, + private_tier: data.private_tier, + active_token_hash: activeTokenHash, + previous_token_hash: current?.active_token_hash && current.active_token_hash !== activeTokenHash + ? current.active_token_hash + : current?.previous_token_hash ?? null, + previous_expires_at: current?.active_token_hash && current.active_token_hash !== activeTokenHash + ? now + data.overlap_ms + : current?.previous_expires_at ?? null, + rotated_by: data.rotated_by, + rotated_at: now, + } - store[key] = rotation - await writeGatewayAuthRotationStore(store) - return rotation + store[key] = rotation + return rotation + }, + ) } async function readPreferenceStore(): Promise { - try { - const raw = await readFile(serverDataJsonPath("notificationPreferences"), "utf8") - const parsed = JSON.parse(raw) as NotificationPreferenceStore - return parsed && typeof parsed === "object" ? parsed : {} - } catch { - return {} - } + return readJsonFile( + serverDataJsonPath("notificationPreferences"), + {}, + ) } -async function writePreferenceStore(data: NotificationPreferenceStore): Promise { - const filePath = serverDataJsonPath("notificationPreferences") - await mkdir(path.dirname(filePath), { recursive: true }) - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8") -} export async function getNotificationPreferences(wallet: string): Promise { const store = await readPreferenceStore() @@ -184,16 +170,24 @@ export async function saveNotificationPreferences( wallet: string, preferences: Partial>, ): Promise { - const current = await getNotificationPreferences(wallet) - const next: NotificationPreferences = { - enabled: preferences.enabled ?? current.enabled, - types: { ...current.types, ...(preferences.types ?? {}) }, - updated_at: Date.now(), - } - const store = await readPreferenceStore() - store[wallet] = next - await writePreferenceStore(store) - return next + // Read the merged defaults inside the lock so two concurrent preference + // updates cannot both start from the same base and lose one field. + return updateStore( + "notificationPreferences", + (store) => { + const current = { + ...DEFAULT_NOTIFICATION_PREFERENCES, + ...(store[wallet] ?? {}), + } + const next: NotificationPreferences = { + enabled: preferences.enabled ?? current.enabled, + types: { ...current.types, ...(preferences.types ?? {}) }, + updated_at: Date.now(), + } + store[wallet] = next + return next + }, + ) } export async function shouldStoreNotification(wallet: string, type: NotificationType): Promise { @@ -204,17 +198,7 @@ export async function shouldStoreNotification(wallet: string, type: Notification } async function readStore(): Promise { - try { - return JSON.parse(await readFile(serverDataJsonPath("notifications"), "utf8")) as NotificationStore - } catch { - return {} - } -} - -async function writeStore(data: NotificationStore): Promise { - const filePath = serverDataJsonPath("notifications") - await mkdir(path.dirname(filePath), { recursive: true }) - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8") + return readJsonFile(serverDataJsonPath("notifications"), {}) } export async function createNotification( @@ -224,8 +208,6 @@ export async function createNotification( ): Promise { if (!(await shouldStoreNotification(wallet, type))) return - const store = await readStore() - const list = store[wallet] ?? [] const notif: Notification = { id: randomUUID(), wallet, @@ -234,10 +216,11 @@ export async function createNotification( created_at: Date.now(), data, } - // Prepend newest first; cap at MAX_PER_WALLET - const updated = [notif, ...list].slice(0, MAX_PER_WALLET) - store[wallet] = updated - await writeStore(store) + await updateStore("notifications", (store) => { + const list = store[wallet] ?? [] + // Prepend newest first; cap at MAX_PER_WALLET + store[wallet] = [notif, ...list].slice(0, MAX_PER_WALLET) + }) } export async function createNotificationBatch( @@ -247,41 +230,42 @@ export async function createNotificationBatch( ): Promise<{ succeeded: number; failed: number }> { if (!wallets.length) return { succeeded: 0, failed: 0 } - const store = await readStore() - let succeeded = 0 - let failed = 0 const now = Date.now() const notificationId = randomUUID() + // Preference checks read a different store, so resolve them before taking the + // notifications lock instead of holding it across the awaits. + const eligible: string[] = [] + let failed = 0 for (const wallet of wallets) { try { - if (!(await shouldStoreNotification(wallet, type))) { - failed++ - continue - } - - const list = store[wallet] ?? [] - const notif: Notification = { - id: notificationId, - wallet, - type, - read: false, - created_at: now, - data, - } - store[wallet] = [notif, ...list].slice(0, MAX_PER_WALLET) - succeeded++ + if (await shouldStoreNotification(wallet, type)) eligible.push(wallet) + else failed++ } catch { failed++ } } - // Batch write - if (succeeded > 0) { - await writeStore(store) - } - - return { succeeded, failed } + if (eligible.length === 0) return { succeeded: 0, failed } + + return updateStore( + "notifications", + (store) => { + for (const wallet of eligible) { + const notif: Notification = { + id: notificationId, + wallet, + type, + read: false, + created_at: now, + data, + } + const list = store[wallet] ?? [] + store[wallet] = [notif, ...list].slice(0, MAX_PER_WALLET) + } + return { succeeded: eligible.length, failed } + }, + ) } export async function getNotifications(wallet: string, limit = 30): Promise { @@ -290,19 +274,19 @@ export async function getNotifications(wallet: string, limit = 30): Promise { - const store = await readStore() - const list = store[wallet] - if (!list) return - store[wallet] = list.map((n) => (n.id === notificationId ? { ...n, read: true } : n)) - await writeStore(store) + await updateStore("notifications", (store) => { + const list = store[wallet] + if (!list) return + store[wallet] = list.map((n) => (n.id === notificationId ? { ...n, read: true } : n)) + }) } export async function markAllRead(wallet: string): Promise { - const store = await readStore() - const list = store[wallet] - if (!list) return - store[wallet] = list.map((n) => ({ ...n, read: true })) - await writeStore(store) + await updateStore("notifications", (store) => { + const list = store[wallet] + if (!list) return + store[wallet] = list.map((n) => ({ ...n, read: true })) + }) } export async function getUnreadCount(wallet: string): Promise { @@ -387,20 +371,11 @@ export class FaucetFunnelError extends Error { } async function readFunnelStore(): Promise { - try { - const raw = await readFile(serverDataJsonPath("faucetFunnelEvents"), "utf8") - const parsed = JSON.parse(raw) as FaucetFunnelStore - return parsed && Array.isArray(parsed.events) ? parsed : { events: [] } - } catch { - return { events: [] } - } + return readJsonFile(serverDataJsonPath("faucetFunnelEvents"), { + events: [], + }) } -async function writeFunnelStore(data: FaucetFunnelStore): Promise { - const filePath = serverDataJsonPath("faucetFunnelEvents") - await mkdir(path.dirname(filePath), { recursive: true }) - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8") -} export async function recordFaucetFunnelEvent( input: unknown, @@ -422,10 +397,17 @@ export async function recordFaucetFunnelEvent( ts: parsed.data.ts ?? opts.now ?? Date.now(), reason: parsed.data.reason ?? null, } - const store = await readFunnelStore() - store.events = [...store.events, event].slice(-MAX_FUNNEL_EVENTS) - await writeFunnelStore(store) - return event + return updateStoreWithReader( + "faucetFunnelEvents", + (raw) => { + const parsed = raw as FaucetFunnelStore | undefined + return parsed && Array.isArray(parsed.events) ? parsed : { events: [] } + }, + (store) => { + store.events = [...store.events, event].slice(-MAX_FUNNEL_EVENTS) + return event + }, + ) } export type ClaimFunnelStageMetric = { diff --git a/lib/profile-store.ts b/lib/profile-store.ts index dc34a64..25dc172 100644 --- a/lib/profile-store.ts +++ b/lib/profile-store.ts @@ -1,6 +1,5 @@ -import { mkdir, readFile, writeFile } from "node:fs/promises"; -import path from "node:path"; import { z } from "zod"; +import { readJsonFile, updateJsonFile, updateStore } from "@/lib/json-store"; import { serverDataJsonPath } from "@/lib/server-data-paths"; import { isFeatureEnabled } from "@/lib/feature-flags"; @@ -18,18 +17,7 @@ export type ProfileData = { type ProfileStore = Record; async function readStore(): Promise { - try { - const raw = await readFile(serverDataJsonPath("profileSocials"), "utf8"); - return JSON.parse(raw) as ProfileStore; - } catch { - return {}; - } -} - -async function writeStore(data: ProfileStore): Promise { - const filePath = serverDataJsonPath("profileSocials"); - await mkdir(path.dirname(filePath), { recursive: true }); - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8"); + return readJsonFile(serverDataJsonPath("profileSocials"), {}); } export async function getProfile(wallet: string): Promise { @@ -41,13 +29,13 @@ export async function saveProfile( wallet: string, data: Omit, ): Promise { - const store = await readStore(); const entry: ProfileData = { ...data, updated_at: Date.now(), }; - store[wallet] = entry; - await writeStore(store); + await updateStore("profileSocials", (store) => { + store[wallet] = entry; + }); return entry; } @@ -217,19 +205,10 @@ function validateProfileHandle(handle: string): string | null { } async function readArtistAliasStore(): Promise { - try { - const raw = await readFile(serverDataJsonPath("artistProfiles"), "utf8"); - const parsed = JSON.parse(raw) as ArtistAliasStore; - return parsed && typeof parsed === "object" ? parsed : {}; - } catch { - return {}; - } -} - -async function writeArtistAliasStore(data: ArtistAliasStore): Promise { - const filePath = serverDataJsonPath("artistProfiles"); - await mkdir(path.dirname(filePath), { recursive: true }); - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8"); + return readJsonFile( + serverDataJsonPath("artistProfiles"), + {}, + ); } export async function resolveProfileHandle( @@ -308,13 +287,23 @@ export async function saveProfileHandle( }; } - const store = await readArtistAliasStore(); - const taken = Object.entries(store).find( - ([wallet, profile]) => - wallet !== walletAddress && - normalizeProfileHandle(profile.alias) === normalized, + // Uniqueness is checked inside the store lock. Reading first and writing + // after let two wallets concurrently claim the same handle. + const claim = await updateStore( + "artistProfiles", + (store) => { + const taken = Object.entries(store).find( + ([wallet, profile]) => + wallet !== walletAddress && + normalizeProfileHandle(profile.alias) === normalized, + ); + if (taken) return "taken"; + store[walletAddress] = { alias: normalized, updatedAt: Date.now() }; + return "ok"; + }, ); - if (taken) { + + if (claim === "taken") { return { ok: false, error: "Handle is already linked to another wallet", @@ -324,8 +313,6 @@ export async function saveProfileHandle( } const updatedAt = Date.now(); - store[walletAddress] = { alias: normalized, updatedAt }; - await writeArtistAliasStore(store); return { ok: true, walletAddress, @@ -809,29 +796,26 @@ export async function recordTrendingSignal( bucketSizeMs: number = 3600000, ): Promise { if (!isPhase87Enabled()) return; - const store = await readJson( - serverDataJsonPath("trendingSignals"), - ); const now = Date.now(); const bucketStart = Math.floor(now / bucketSizeMs) * bucketSizeMs; const bucketEnd = bucketStart + bucketSizeMs; - const signal = store[wallet] ?? { wallet, score: 0, buckets: [] }; - let bucket = signal.buckets.find((b) => b.start === bucketStart); + await updateStore("trendingSignals", (store) => { + const signal = store[wallet] ?? { wallet, score: 0, buckets: [] }; + let bucket = signal.buckets.find((b) => b.start === bucketStart); - if (!bucket) { - bucket = { start: bucketStart, end: bucketEnd, views: 0, engagement: 0 }; - signal.buckets.push(bucket); - } - - bucket.views++; - signal.score = signal.buckets.reduce( - (sum, b) => sum + b.views + b.engagement, - 0, - ); - store[wallet] = signal; + if (!bucket) { + bucket = { start: bucketStart, end: bucketEnd, views: 0, engagement: 0 }; + signal.buckets.push(bucket); + } - await writeJson(serverDataJsonPath("trendingSignals"), store); + bucket.views++; + signal.score = signal.buckets.reduce( + (sum, b) => sum + b.views + b.engagement, + 0, + ); + store[wallet] = signal; + }); } export async function getTrendingSignals( @@ -1089,19 +1073,7 @@ export function auditFaucetRateLimitWiring(): { ok: boolean; note: string } { } async function readJson(filePath: string): Promise { - try { - return JSON.parse(await readFile(filePath, "utf8")) as T; - } catch { - return {} as T; - } -} - -async function writeJson( - filePath: string, - data: T, -): Promise { - await mkdir(path.dirname(filePath), { recursive: true }); - await writeFile(filePath, JSON.stringify(data, null, 2), "utf8"); + return readJsonFile(filePath, {} as T); } // ── Issues #65 / #66 (phase-137): structured error taxonomy ───────────────────