Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion PROJECT_ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ Owns:
- Signals, replies, edit history and reactions live in SQLite (`lib/sqlite-db.ts`, WAL mode), not the JSON sidecar.
- x402 endpoints and payment verification support.
- Concurrency safety for the signal store: version-guarded writes (`WHERE id = ? AND version = ?`) with bounded retry, plus `BEGIN IMMEDIATE` transactions for parent-version checks. Holds across processes, so concurrent Vercel instances no longer interleave. See `docs/TECHNICAL.md` 5.3.
- Collaborative signal lore drafting (phase-141, issue #207): a Yjs CRDT draft in SQLite (`signal_crdt_docs` / `signal_crdt_updates`) that merges concurrent edits automatically and never writes the `signals` row. The merged draft is promoted to committed lore only through `PUT /api/signals/[id]` with `from_draft: true`, which reuses the same `If-Match` CAS and `signal_versions` snapshot as a manual edit. Sync is HTTP state-vector exchange, not WebSocket, because the app deploys to Vercel serverless. See `docs/TECHNICAL.md` 5.3.
- Other JSON-backed stores (`follow`, `profile`, `market`, `notification`, `achievement`, `narrative-world`) still do unguarded read-modify-write on their sidecar files. Single-writer is safe; concurrent writers lose updates.

Must not own:
Expand Down Expand Up @@ -103,7 +104,14 @@ Owns:
mutations may echo it back as `If-Match` (or `parent_version` for appends) and
the server rejects a stale value with `409` rather than silently overwriting a
concurrent writer. `409` is logged under the `signals.version_conflict` event
so conflict rates are observable.
and counted in `signal_version_conflicts` so conflict rates are observable.
- **Optimistic concurrency guards commits, it does not refuse collaboration.**
Where many people edit one record at once, `409` on every save is a correct
answer to the wrong question. Signals therefore layer a CRDT *draft* over the
guarded record: concurrent proposals merge automatically, and the merged
result is promoted through the same `If-Match` guard so the authoritative row
keeps a single linear history. The guard is never removed; it moves to the one
place a decision actually has to be made.

## 6) Internationalization architecture

Expand Down Expand Up @@ -184,6 +192,7 @@ is a fixed-size digest of the canonical payload, keeping it under wallet
| `phase-83` | `NEXT_PUBLIC_FEATURE_PHASE_83` / `FEATURE_PHASE_83` | Emoji-reaction aggregation on signals (curated set, toggle per wallet) with a 20/60s per-wallet rate limit | off | Unset var, restart — reactions route returns 404; existing `signal_reactions` rows remain on disk (no migration to undo) |
| `phase-139` | `NEXT_PUBLIC_FEATURE_PHASE_139` / `FEATURE_PHASE_139` | Collection-level offer books aggregated from per-token offers, plus bulk-bid across a collection's listings | off | Unset var, restart — offer-book/bulk-bid route returns 404; per-listing offers (`/api/market/[id]/offers`) are unaffected either way |
| `phase-140` | `NEXT_PUBLIC_FEATURE_PHASE_140` / `FEATURE_PHASE_140` | Royalty enforcement on secondary sales: creator/seller split computed and ledgered at offer-accept time | off | Unset var, restart — listing creation stops accepting `creator_wallet`/`royalty_bps`; offer-accept stops computing a split (100% to seller, pre-140 behavior); existing `royalty_payouts` rows are historical record |
| `phase-141` | `NEXT_PUBLIC_FEATURE_PHASE_141` / `FEATURE_PHASE_141` | Yjs CRDT collaborative draft for signal lore — concurrent edits merge automatically instead of clobbering, promoted to committed lore via `PUT` (`from_draft: true`) under the normal `If-Match` guard | off | Unset var, restart — the CRDT routes return 404 and the draft panel is hidden; `PUT /api/signals/[id]` still works on its own as an `If-Match`-guarded full replacement; existing `signal_crdt_*` rows are inert scratch state (no migration to undo) |

Flags are read via `lib/feature-flags.ts:isFeatureEnabled`. Client flags use `NEXT_PUBLIC_*`, server also accepts `FEATURE_*`. Zero regression when off.

Expand Down
164 changes: 164 additions & 0 deletions app/api/signals/[id]/crdt/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
import { NextRequest } from "next/server"
import { StrKey } from "@stellar/stellar-sdk"
import { createApiRequestContext } from "@/lib/api-observability"
import { getSignal } from "@/lib/signal-store"
import {
isSignalCrdtEnabled,
flag141RollbackNote,
mergeSignalLoreUpdate,
readSignalLoreDraft,
SignalCrdtError,
} from "@/lib/signal-crdt-store"

export const runtime = "nodejs"
export const dynamic = "force-dynamic"

/**
* Issue #207 (phase-141): the Yjs sync surface for a signal's collaborative
* lore draft.
*
* This is the same two-step sync protocol `y-websocket` speaks — client sends
* its state vector, server replies with the operations it is missing; client
* sends an update, server folds it in and replies with its own state vector —
* carried over plain HTTP instead of a WebSocket. The transport differs because
* this app deploys to Vercel serverless, where a request has a bounded lifetime
* and no upgrade handshake is possible; the merge semantics, which are what
* actually prevent lost writes, are Yjs' either way. The measurements behind
* that trade-off are in `docs/spikes/207-signal-crdt-benchmark.md`.
*
* `GET` `?state_vector=<base64>` → the operations the caller is missing.
* `POST` `{ update, wallet }` → fold an update, return the merged draft.
*/
const CRDT_ERROR_STATUS: Record<SignalCrdtError["code"], number> = {
FLAG_DISABLED: 404,
NOT_FOUND: 404,
VALIDATION_FAILED: 400,
}

export async function GET(
request: NextRequest,
{ params }: { params: Promise<{ id: string }> },
) {
const api = createApiRequestContext(request, "/api/signals/[id]/crdt")
const { id } = await params

if (!isSignalCrdtEnabled()) {
return api.json(
{ error: "Collaborative lore drafting disabled (phase-141 flag off)", rollback: flag141RollbackNote() },
{ status: 404, event: "signals.crdt.disabled" },
)
}

const rawStateVector = request.nextUrl.searchParams.get("state_vector")?.trim() || undefined

try {
const signal = await getSignal(id)
if (!signal) {
return api.json(
{ error: "Signal not found" },
{ status: 404, event: "signals.crdt.signal_missing", metadata: { signal_id: id } },
)
}
const draft = await readSignalLoreDraft(id, rawStateVector)
return api.json(
{ ...draft, version: signal.version },
{
event: "signals.crdt.synced",
metadata: { signal_id: id, update_count: draft.updateCount, contributors: draft.contributors.length },
},
)
} catch (error) {
return crdtFailure(error, api, id)
}
}

type SyncBody = {
update?: unknown
wallet?: unknown
/** Optional: narrows the reply to only the operations this client lacks. */
state_vector?: unknown
}

export async function POST(
request: NextRequest,
{ params }: { params: Promise<{ id: string }> },
) {
const api = createApiRequestContext(request, "/api/signals/[id]/crdt")
const { id } = await params

if (!isSignalCrdtEnabled()) {
return api.json(
{ error: "Collaborative lore drafting disabled (phase-141 flag off)", rollback: flag141RollbackNote() },
{ status: 404, event: "signals.crdt.disabled" },
)
}

let body: SyncBody
try {
body = (await request.json()) as SyncBody
} catch {
return api.json({ error: "Invalid JSON" }, { status: 400, event: "signals.crdt.invalid_json" })
}

if (typeof body.wallet !== "string" || !StrKey.isValidEd25519PublicKey(body.wallet)) {
return api.json(
{ error: "Invalid wallet address" },
{ status: 400, event: "signals.crdt.validation_failed", metadata: { reason: "wallet" } },
)
}
if (typeof body.update !== "string" || body.update.length === 0) {
return api.json(
{ error: "update required" },
{ status: 400, event: "signals.crdt.validation_failed", metadata: { reason: "update" } },
)
}
if (body.state_vector != null && (typeof body.state_vector !== "string" || body.state_vector.length === 0)) {
return api.json(
{ error: "state_vector must be a base64 string when present" },
{ status: 400, event: "signals.crdt.validation_failed", metadata: { reason: "state_vector" } },
)
}
const stateVector = typeof body.state_vector === "string" ? body.state_vector : undefined

try {
const signal = await getSignal(id)
if (!signal) {
return api.json(
{ error: "Signal not found" },
{ status: 404, event: "signals.crdt.signal_missing", metadata: { signal_id: id } },
)
}

const { draft, concurrent } = await mergeSignalLoreUpdate(id, body.update, body.wallet, {
sinceStateVector: stateVector,
})

return api.json(
{ ...draft, version: signal.version, concurrent },
{
event: "signals.crdt.merged",
metadata: { signal_id: id, concurrent, update_count: draft.updateCount },
},
)
} catch (error) {
return crdtFailure(error, api, id)
}
}

function crdtFailure(
error: unknown,
api: ReturnType<typeof createApiRequestContext>,
signalId: string,
) {
if (error instanceof SignalCrdtError) {
return api.json(
{ error: error.message, code: error.code },
{
status: CRDT_ERROR_STATUS[error.code],
event: error.code === "NOT_FOUND" ? "signals.crdt.signal_missing" : "signals.crdt.rejected",
metadata: { signal_id: signalId, reason: error.code },
},
)
}
return api.errorJson(error, 500, "signals.crdt.failed")
}
Loading