From 4d879997995dc71dbc33751938e38f256dfc5fca Mon Sep 17 00:00:00 2001 From: Mikey-222 Date: Mon, 28 Sep 2026 09:02:05 +0100 Subject: [PATCH] fix: extend the sidecar write lock to the six stores #361 left unguarded #361 fixed lib/signal-store.ts and its commit message noted that the other six JSON stores "share the same unsafe pattern and are untouched here". This does that: each one hand-rolled a readFile -> JSON.parse -> mutate -> writeFile pair with no coordination, so concurrent read-modify-write cycles raced and the second writeFile discarded the first. Measured on this branch before the change, 50 concurrent writers: follows 1/50 persisted profile saves 1/50 persisted checkAndUnlock counters 1/50 accumulated world save versions 1 (every writer read version 0) profile view counters 1/50 counted unlockAchievement reported success 10/10 for the same unlock Adds lib/json-store.ts, giving the same three guarantees #361 applied inline to signals, and routes the remaining stores through it: serialized mutation, atomic temp-then-rename replacement, and a JSON parse failure that throws instead of degrading to {} (degrading meant the next write wiped the store). Also fixes a bug needing no concurrency at all: checkAndUnlock read the store, let tryUnlock call unlockAchievement (its own read-modify-write), then wrote its stale snapshot back. It returned the unlock to its caller and fired an "achievement unlocked" notification for an achievement that was never persisted. The counter update and every unlock it triggers now happen in one locked mutation. Because the lock is not reentrant, unlocks are applied to the in-flight store rather than re-entering through unlockAchievement. Two further read-then-write races are now decided under the lock: claimProfileHandle's alias uniqueness check (two wallets could claim one handle) and setNotificationPreferences' merge (concurrent updates each started from the same base and lost a field). Also fixes a pre-existing runtime bug found while testing: markNarrativeRead and getReaderProgress passed "readerProgress" to serverDataJsonPath, but the key is "worldReaderProgress", so FILES[key] was undefined and path.join(root, undefined) threw. Reader progress was entirely non-functional. @ts-nocheck on that file is why tsc never caught it. Scope: the six stores #361 named. market-store's listings and offers are already SQLite, so only marketProfileViews and blockList are touched here. The route-local stores (faucet claims, classic-liq claims, artist profiles, nft listings) still have the original pattern and are not addressed here. Tests: lib/__tests__/json-store-concurrency.test.ts, 22 vitest cases, 12 of which fail against the unpatched stores. Adds no new type errors (163 vs the 164 baseline, all pre-existing on main). Not addressed, and reported separately: - The lock is per-process; Vercel instances each hold their own os.tmpdir() copy, so cross-instance state still diverges. - main is broken independently of this change: 163 type errors, plus lib/signal-store.ts has a duplicate `let items` declaration that is a hard parse error, so that module and both of its test files cannot load. - 40 of the 43 test files on main use node:test, which `vitest run` cannot collect, so `npm test` reports "No test suite found" for all of them. --- docs/TECHNICAL.md | 38 +++ lib/__tests__/json-store-concurrency.test.ts | 278 +++++++++++++++++++ lib/achievement-store.ts | 193 +++++++------ lib/follow-store.ts | 41 +-- lib/json-store.ts | 167 +++++++++++ lib/market-store.ts | 145 +++++----- lib/narrative-world-store.ts | 161 ++++++----- lib/notification-store.ts | 230 +++++++-------- lib/profile-store.ts | 108 +++---- 9 files changed, 910 insertions(+), 451 deletions(-) create mode 100644 lib/__tests__/json-store-concurrency.test.ts create mode 100644 lib/json-store.ts diff --git a/docs/TECHNICAL.md b/docs/TECHNICAL.md index aa025ff..71cf7e3 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) | @@ -214,6 +215,39 @@ cross-process transactions (Postgres `UPDATE ... WHERE version = $expected`). This is deliberate: it is a bounded, documented limit rather than a CRDT layer that would not have fixed the underlying file race. +### 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 @@ -317,6 +351,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 ───────────────────