diff --git a/.straymark/07-ai-audit/agent-logs/AILOG-2026-07-16-003-charter-14-durabilidad-relay.md b/.straymark/07-ai-audit/agent-logs/AILOG-2026-07-16-003-charter-14-durabilidad-relay.md new file mode 100644 index 0000000..fdfe390 --- /dev/null +++ b/.straymark/07-ai-audit/agent-logs/AILOG-2026-07-16-003-charter-14-durabilidad-relay.md @@ -0,0 +1,100 @@ +--- +id: AILOG-2026-07-16-003 +title: "CHARTER-14: durabilidad del relay — persist-before-broadcast + fsync + cobertura de carga del relay" +status: accepted +created: 2026-07-16 +agent: claude-opus-4-8 +confidence: high +review_required: true +risk_level: medium +eu_ai_act_risk: not_applicable +nist_genai_risks: [] +iso_42001_clause: [] +observability_scope: none +tags: [relay, durability, persist-before-broadcast, fsync, load-test, follow-ups, m3] +related: [AIDEC-2026-07-16-001, AILOG-2026-07-16-002] +originating_charter: CHARTER-14-durabilidad-relay +--- + +# AILOG: CHARTER-14 — durabilidad del relay (FU-010) + +## Summary + +Despacho de FU-010, último Charter del vaciado del backlog antes del publish (T060). Invierte el orden +del relay a **persist-before-broadcast** (default), añade **fsync** al store de referencia y la +**cobertura de carga del relay** que no existía. Con esto el backlog baja de 2 open a **1** (FU-015, +bloqueado upstream) — el vaciado está completo salvo lo que depende de terceros. + +Decisiones de diseño en AIDEC-2026-07-16-001 (supersede de AIDEC-2026-07-13-001 §5). Tres las fijó el +operador ex-ante (fsync en scope, default persist-first, cobertura de carga); el AILOG registra la +ejecución y los hallazgos. + +## Actions Performed + +1. **`WeftServerOptions.Durability`** (enum `DurabilityMode`, default `PersistThenBroadcast`). +2. **`DocumentSession.ApplyAndCaptureDeltaAsync`**: aplica en el turno del actor y **devuelve el delta** + (race-free vía valor de retorno, frente a conexiones concurrentes del mismo hub). Refinamiento sobre + el plan: el hub deja de usar el evento `UpdateApplied` para el broadcast (ver AIDEC Decisión 4). +3. **`DocumentHub`**: broadcast explícito ordenado por modo; en `PersistThenBroadcast`, fallo de append + → `DisconnectAll()` + relanza (el update quedó en el doc vivo sin difundir → cerrar el documento + fuerza el re-sync autoritativo). +4. **`WeftConnection.ApplyOrCloseAsync`**: fallo del store → cierre **1011** (antes escapaba de + `RunAsync`, que solo capturaba OCE/WebSocketException). Cancelación y socket cortado se dejan propagar. +5. **`FileSystemDocumentStore`**: `Flush(flushToDisk: true)` + fsync del directorio (POSIX, vía + `RandomAccess.FlushToDisk`; omitido en Windows). +6. **`WeftServer`**: pasa `_options.Durability` al `new DocumentHub(...)`. +7. **Cobertura de carga del relay** (`Weft.LoadTest`, modo `--relay`): editor→observador vía TestServer + real contra `FileSystemDocumentStore` con fsync, en ambos modos, con p50/p99. Nuevas refs: + `Weft.Server` + `Microsoft.AspNetCore.TestHost`. +8. **Tests de relay** (`RelayTests`, +4): fallo de append no observado + cierre 1011; reconexión + resincroniza; orden en ambos modos (testigo con store armado). `ControllableAppendStore` con «armado» + —los appends del handshake pasan, el test arma el fallo/bloqueo justo antes de la edición— porque + contar appends por número es frágil (cada conexión hace un append en el SyncStep2). +9. **AIDEC** (supersede §5) + este AILOG + cierre de FU-010 + `recount`. + +## Risk + +Riesgos del Charter (R1–R7) y su desenlace: + +- **R1 (fsync no basta si el store del consumidor no es durable)** — mitigado y documentado: la garantía + es «ningún par observa lo no aceptado por el store»; la durabilidad física del aceptado la aporta el + adaptador. El fsync cubre el store de referencia (filesystem). +- **R2 (divergencia permanente si el fallo de append no cierra las conexiones)** — mitigado, es + obligatorio: `DisconnectAll` + 1011, con test dedicado de fallo inyectado determinista. +- **R3 (regresión de latencia)** — **medido, y no se materializó**: p50/p99 idénticos entre modos + (2.1ms), el coste del orden seguro es ~0 en percentiles típicos. Solo el `max` sube (18.6 vs 2.5ms) + por picos de fsync. La carga del relay es la evidencia que respalda el default. +- **R4 (reorden cross-conexión rompe un cliente no-Yjs)** — mitigado: acotado por el pending store de + yrs; convergencia real de 2 clientes `y-websocket` verificada end-to-end con el default nuevo. +- **R5 (blip del store desconecta todo un doc)** — aceptado y documentado: retry/backoff es del + adaptador; los clientes reconectan. +- **R6 (doble ruta duplica el espacio de estados)** — mitigado: ambos modos por la misma clase de tests; + comparten todo salvo el orden de dos líneas. +- **R7 (fsync síncrono bloquea un hilo del pool)** — medido: despreciable en p50/p99 (el append ya está + fuera del turno del actor). Si escalara, `RandomAccess.FlushToDisk` ya se usa para el directorio; el + archivo usa `Flush(flushToDisk:true)`. No hizo falta optimizar. + +Sin R8: no surgió ningún riesgo nuevo no anticipado. + +## Follow-ups + +Ninguno nuevo. FU-010 cerrado. El vaciado del backlog (CHARTER-12/13/14) deja **1 open**: FU-015 +(adopción del fix de R6 vía bump de yrs), bloqueado por el merge upstream de y-crdt#639 — no accionable +por nosotros. + +Nota (no accionable): un falso positivo conocido de `charter drift` con `.csproj`/`.sln` (#354) puede +reaparecer al declarar `Weft.LoadTest.csproj`; el archivo está declarado y su cambio (refs a Weft.Server ++ TestHost) es intencional. + +## Verification + +```bash +cargo build --release --features test-hooks --manifest-path native/Cargo.toml +dotnet test Weft.sln -c Release # 155/155 (Server 70→74) +dotnet test tests/Weft.Server.Tests -c Release # incl. contract suite intacta (IDocumentStore no cambia) +dotnet run --project tests/Weft.LoadTest -c Release -- --relay --edits 200 # PASS; p50/p99 idénticos entre modos +# Convergencia real con el default nuevo (reorden no rompe Yjs): +dotnet run --project samples/Weft.Sample.Server -c Release & +cd samples/tiptap-client && npm run check # "Hello from A. And B too." +straymark validate --include-charters +``` diff --git a/.straymark/07-ai-audit/decisions/AIDEC-2026-07-16-001-charter-14-durabilidad-relay-persist-before-broadcast.md b/.straymark/07-ai-audit/decisions/AIDEC-2026-07-16-001-charter-14-durabilidad-relay-persist-before-broadcast.md new file mode 100644 index 0000000..50c32b3 --- /dev/null +++ b/.straymark/07-ai-audit/decisions/AIDEC-2026-07-16-001-charter-14-durabilidad-relay-persist-before-broadcast.md @@ -0,0 +1,147 @@ +--- +id: AIDEC-2026-07-16-001 +title: "CHARTER-14: durabilidad del relay — persist-before-broadcast por defecto (supersede de AIDEC-2026-07-13-001 §5)" +status: accepted +created: 2026-07-16 +agent: claude-opus-4-8 +confidence: high +review_required: true +risk_level: medium +eu_ai_act_risk: not_applicable +nist_genai_risks: [] +iso_42001_clause: [] +tags: [relay, durability, persist-before-broadcast, fsync, backpressure, ordering] +related: [AILOG-2026-07-16-003, AIDEC-2026-07-13-001] +originating_charter: CHARTER-14-durabilidad-relay +--- + +# AIDEC: durabilidad del relay — persist-before-broadcast por defecto + +> Registra las decisiones de diseño de CHARTER-14 (FU-010). **Supersede la decisión §5 de +> AIDEC-2026-07-13-001** (CHARTER-05), que fijó broadcast-then-persist como comportamiento del relay con +> la auto-sanación CRDT como red. Aquella decisión sigue siendo válida como *modo opcional*; deja de ser +> el default. + +## Context + +El relay difundía cada update a los pares **antes** de persistirlo (`DocumentSession.UpdateApplied` se +dispara dentro del turno del actor, antes de `AppendUpdateAsync`). AIDEC-2026-07-13-001 §5 lo justificó: +en single-node, un append fallido + crash se recupera por re-sync en la reconexión. La auditoría de +CHARTER-05 (F3) lo marcó como follow-up (FU-010) para cuando se requieran garantías de durabilidad duras. +El operador decidió implementarlo como último Charter del vaciado del backlog, con tres sub-decisiones +tomadas ex-ante (fsync en scope, default persist-first, cobertura de carga). Este AIDEC registra el +*cómo* y el *porqué* de la forma concreta. + +--- + +## Decisión 1 — Default = PersistThenBroadcast (no el comportamiento heredado) + +### Problem + +¿El default debe ser el orden seguro (ningún par ve lo no persistido) o el heredado (menor latencia, +auto-sanación como red)? Cambiar el default es un breaking change si el repo ya es público. + +### Decision + +**Default `PersistThenBroadcast`**, con `BroadcastThenPersist` como válvula de escape opt-in +(`WeftServerOptions.Durability`). + +### Rationale + +- El coste es **una latencia que el emisor ya paga**: `AppendUpdateAsync` ya se espera en el receive + loop antes de leer el siguiente frame; invertir el orden no añade I/O a ningún hot path, solo mueve el + broadcast a después del append. **Medido** (carga del relay, ambos modos, `FileSystemDocumentStore` + con fsync): p50 y p99 **idénticos** (2.1ms), el coste del orden seguro es ~0 en percentiles típicos; + solo el `max` sube (18.6ms vs 2.5ms) por picos ocasionales de fsync. La premisa del diseño se confirma + con datos, no con criterio. +- El semantic seguro debe ser el default. La garantía que compra —*ningún par observa estado que el + servidor no haya aceptado de forma durable*— es la que un operador espera por defecto de un relay. +- El repo aún no es público (FU-014): cambiar el default ahora es gratis; hacerlo tras el publish sería + breaking. Es el momento correcto. + +--- + +## Decisión 2 — Fallo de append ⇒ cerrar TODAS las conexiones del documento (1011) + +### Problem + +Con persist-first, un append fallido deja el update en el doc vivo pero **nunca difundido**. Los pares +quedan callados y desactualizados; la próxima edición difunde solo su delta, jamás el que faltó → +divergencia permanente. ¿Cómo se recupera? + +### Decision + +En `PersistThenBroadcast`, un fallo de `AppendUpdateAsync` dispara `DocumentHub.DisconnectAll()` y la +conexión emisora cierra con **1011** (InternalError). Todos reconectan, mandan SyncStep1, y el servidor +—autoritativo, que sí tiene el update en el doc vivo— reenvía el estado. Convergencia recuperada. + +### Rationale + +Es el **único modo de equivocarse** en persist-first, así que la remediación es obligatoria, no +opcional. Cerrar el documento entero es contundente pero correcto: fuerza el re-sync autoritativo desde +la única fuente que tiene el estado completo. Retry/backoff es responsabilidad del adaptador del store, +no de Weft; los clientes reconectan solos. Cubierto por un test con fallo de append inyectado +determinista (sin matar procesos): se aserta que ningún par recibió el delta, que la conexión cerró +1011, y que un cliente que reconecta resincroniza. + +--- + +## Decisión 3 — fsync de archivo + directorio en FileSystemDocumentStore + +### Problem + +`FlushAsync` solo vacía al page cache del SO; `File.Move` no es durable sin fsync del directorio. Sin +fsync, persist-first protege solo contra crash de **proceso**, no de máquina — y el trigger de FU-010 +dice «SLA de no-pérdida». + +### Decision + +`fs.Flush(flushToDisk: true)` (síncrono; no hay fsync async en .NET) + fsync del directorio contenedor +en POSIX vía `RandomAccess.FlushToDisk`. En Windows se omite el fsync de directorio (no hay handle +equivalente; NTFS ordena la metadata de forma que el rename no precede al contenido ya sincronizado). + +### Rationale + +La garantía de orden («ningún par observa lo no aceptado») solo vale si «aceptado» significa «durable en +disco», no «en page cache». El fsync síncrono bloquea un hilo del pool durante la operación, pero el +append ya está fuera del turno del actor (no bloquea lecturas del doc). El coste se midió (ver Decisión +1): despreciable en p50/p99. **Límite honesto**: esto cubre el store de referencia (filesystem); Redis, +EFCore e InMemory tienen su propia semántica de durabilidad, ajena a Weft — la durabilidad física del +«aceptado» la aporta el adaptador de cada backend. + +--- + +## Decisión 4 — Broadcast explícito, no vía el evento (refinamiento sobre el plan) + +### Problem + +El plan proponía suscribir/no-suscribir `UpdateApplied` según el modo. Pero varias conexiones del mismo +hub pueden llamar `ApplyAndPersistAsync` concurrentemente (distintos receive loops), lo que descarta +capturar el delta en un campo compartido (carrera). + +### Decision + +El hub **deja de usar el evento** para el broadcast del relay. `DocumentSession.ApplyAndCaptureDeltaAsync` +aplica y **devuelve el delta como valor de retorno del turno** (race-free: cada llamada tiene su propio +retorno). `ApplyAndPersistAsync` ordena append/broadcast según el modo. El evento `UpdateApplied` se +conserva para otros consumidores; el relay ya no depende de él. + +### Rationale + +Más limpio que el condicional de suscripción: el modo solo intercambia dos líneas (append↔broadcast) y se +elimina la doble computación de delta (el actor ya no la computa si nadie está suscrito). Race-free por +construcción. El comportamiento de `BroadcastThenPersist` se preserva fielmente (difunde antes del +append), verificado con un testigo de orden por modo. + +--- + +## Consecuencias + +- `WeftServerOptions.Durability` (default `PersistThenBroadcast`); `DocumentSession.ApplyAndCaptureDeltaAsync`; + `DocumentHub` con broadcast explícito + `DisconnectAll` en fallo; `WeftConnection` cierra 1011; + `FileSystemDocumentStore` con fsync; modo `--relay` en `Weft.LoadTest`. +- **No** se cablea ack de aplicación al emisor: y-protocols no lo tiene para `Update`, y la garantía es + sobre lo que observan los pares, no sobre lo que sabe el emisor. +- Reorden de broadcast cross-conexión aceptado (acotado por el pending store de yrs); ruta de escalada + (pump por-hub) registrada, no construida. +- El contrato de `IDocumentStore` **no cambia**: la contract suite (4 adaptadores) sigue verde sin tocar. diff --git a/.straymark/charters/14-durabilidad-relay.md b/.straymark/charters/14-durabilidad-relay.md new file mode 100644 index 0000000..8152a3b --- /dev/null +++ b/.straymark/charters/14-durabilidad-relay.md @@ -0,0 +1,193 @@ +--- +charter_id: CHARTER-14-durabilidad-relay +status: in-progress +effort_estimate: L +trigger: "El operador decide implementar FU-010 como último Charter del vaciado del backlog antes del publish (T060), con tres decisiones tomadas: fsync EN scope, default PersistThenBroadcast, y cobertura de carga del relay (hoy inexistente). Al cerrar, el backlog queda en 1 open (FU-015, bloqueado upstream)." +originating_spec: specs/001-weft-crdt-versioning/spec.md +work_verb: implement +design_provenance: new +--- + +# Charter: Durabilidad del relay — persist-before-broadcast + fsync + +> **Status (mirrored from frontmatter — source of truth is above):** in-progress. Effort: L. +> +> **Origin:** FU-010 (AIDEC-2026-07-13-001 §5, auditoría de CHARTER-05 F3). El relay hace hoy +> **broadcast-then-persist**; este Charter lo invierte, endurece la durabilidad del store con fsync, y +> añade la cobertura de carga del relay que hoy no existe. Tercero y último del vaciado del backlog. + +## Context + +El relay difunde un update a los pares **antes** de persistirlo. El mecanismo real (verificado en HEAD) +es sutil: `DocumentHub.ApplyAndPersistAsync` hace `await Session.ApplyUpdateAsync` y luego +`await _store.AppendUpdateAsync`, pero el **broadcast no está ahí** — lo dispara el evento +`DocumentSession.UpdateApplied` **dentro** del turno del actor durante el apply, es decir *antes* de que +`AppendUpdateAsync` llegue a ejecutarse. Un fallo del append tras el broadcast deja a los pares con un +update que el store no tiene; en v1 single-node la auto-sanación CRDT (re-sync en reconexión) lo +recupera, y así lo documentó AIDEC-2026-07-13-001 §5 como decisión consciente. + +**La premisa del follow-up sobre el coste es falsa, y eso simplifica el trabajo.** FU-010 temía que +persist-before-broadcast metiera I/O en el camino caliente del actor. No es así: `AppendUpdateAsync` +**ya se espera** en el receive loop (`WeftConnection.DispatchAsync`, awaited secuencialmente), así que +cada conexión ya paga la latencia del append antes de leer su siguiente frame. Invertir el orden **no +añade I/O a ningún hot path**: solo mueve el broadcast de dentro del turno a después del append. + +**Refinamiento de diseño sobre el plan.** El plan proponía suscribir/no-suscribir el evento según el +modo. Es más limpio que el hub **deje de usar el evento** y haga el broadcast **explícito**, capturando +el delta como valor de retorno del turno (`DocumentSession.ApplyAndCaptureDeltaAsync`). Así el modo solo +intercambia dos líneas (append↔broadcast) y se elimina la doble computación de delta. Es race-free +porque cada llamada obtiene su propio valor de retorno — necesario, porque **varias conexiones del mismo +hub pueden llamar `ApplyAndPersistAsync` concurrentemente** (distintos receive loops), lo que descarta +capturar el delta en un campo compartido. + +**fsync (decisión del operador: EN scope).** `FileSystemDocumentStore.AtomicWriteAsync` usa +`FlushAsync` → page cache del SO, no disco; `File.Move` tampoco es durable sin fsync del directorio. Sin +fsync, persist-before-broadcast protege solo contra crash de **proceso**, no de máquina — y el trigger +de FU-010 dice «SLA de no-pérdida». Se añade `Flush(flushToDisk: true)` + fsync del directorio (POSIX). + +**Cobertura de carga (decisión del operador: EN scope).** `tests/Weft.LoadTest` conduce +`DocumentBroker` directamente: **cero** referencias a `DocumentHub`/`WeftServer`/`IDocumentStore` +(verificado). Los 45k ops/20s son ciegos a este cambio; no hay hoy ninguna cobertura de carga del path +de persistencia del relay. Se añade un modo que ejercita el relay real vía `TestServer`. + +## Scope + +**In scope:** + +1. **`WeftServerOptions.DurabilityMode`** (enum `PersistThenBroadcast` | `BroadcastThenPersist`) + + propiedad `Durability`. **Default `PersistThenBroadcast`** (el semantic seguro; el repo aún no es + público, cambiar el default ahora es gratis y luego sería breaking). +2. **`DocumentSession.ApplyAndCaptureDeltaAsync`**: aplica el update en el turno del actor y **devuelve + el delta** (vía la ruta de `ExecuteAsync`, race-free). El evento `UpdateApplied` se conserva para + otros consumidores; el relay deja de depender de él. +3. **`DocumentHub`**: deja de suscribir `UpdateApplied`; `ApplyAndPersistAsync` captura el delta, + ordena append/broadcast según el modo, y difunde explícito (eco al emisor conservado, `exclude: + null`). En `PersistThenBroadcast`, **un fallo de append → `DisconnectAll()` + relanza** (los pares + quedarían callados y desactualizados para siempre si no; reconectan y el servidor autoritativo + reenvía el delta). +4. **`WeftConnection`**: un fallo del store en `ApplyAndPersistAsync` cierra la conexión con **1011** + (InternalError) en vez de escapar de `RunAsync` (hoy solo captura OCE/WebSocketException). +5. **`FileSystemDocumentStore`**: `Flush(flushToDisk: true)` + fsync del directorio tras `File.Move`. +6. **Cobertura de carga del relay**: modo nuevo en `tests/Weft.LoadTest` (o sibling) que conduce N + clientes WebSocket vía `TestServer` contra `FileSystemDocumentStore`, reportando latencia de + broadcast en ambos modos — la evidencia que respalda el default. +7. **AIDEC** que supersede a `AIDEC-2026-07-13-001` §5. + +**Out of scope:** + +- **Ack de aplicación al emisor.** y-protocols no lo tiene para `Update`; el emisor nunca sabe si su + update se persistió, con orden o sin él. La propiedad que se compra es exactamente: *ningún par + observa estado que el servidor no haya aceptado de forma durable*. No se inventa un ack propio. +- **Orden estricto de broadcast cross-conexión** (una «durability pump» por-hub). El reordenamiento + cross-conexión es aceptable: los appends ya son desordenados hoy, yrs bufferea updates causalmente + incompletos en su pending store, y el reload reproduce en orden de append. Se documenta la ruta de + escalada sin construirla. +- **FU-015** — bloqueado por upstream. + +## Files to modify + +| File | Change | +|---|---| +| `src/Weft.Server/WeftServerOptions.cs` | `+ DurabilityMode Durability` (default `PersistThenBroadcast`) + el enum | +| `src/Weft.Core/Concurrency/DocumentSession.cs` | `+ ApplyAndCaptureDeltaAsync` (aplica + devuelve el delta) | +| `src/Weft.Server/DocumentHub.cs` | Broadcast explícito según modo; fallo de append → `DisconnectAll` + relanza; deja de suscribir el evento | +| `src/Weft.Server/WeftServer.cs` | Pasa `options.Durability` al `new DocumentHub(...)` (:102) | +| `src/Weft.Server/WeftConnection.cs` | Fallo del store → cierre 1011 (no escapar de `RunAsync`) | +| `src/Weft.Server/Persistence/FileSystemDocumentStore.cs` | `Flush(flushToDisk: true)` + fsync del directorio | +| `tests/Weft.Server.Tests/RelayTests.cs` | Tests: fallo de append no se observa; cierre 1011; reconexión resincroniza; modo legacy; testigo de orden | +| `tests/Weft.LoadTest/RelayLoad.cs` | New — módulo de carga del relay (N clientes WS vía TestServer, p50/p99 por modo) | +| `tests/Weft.LoadTest/Program.cs` | Cablea el modo `--relay` al inicio | +| `tests/Weft.LoadTest/Weft.LoadTest.csproj` | ProjectReference a `Weft.Server` + `Microsoft.AspNetCore.TestHost` | +| `.straymark/follow-ups-backlog.md` | Cierre de FU-010 | +| `.straymark/07-ai-audit/decisions/AIDEC-2026-07-16-001-durabilidad-relay-persist-before-broadcast.md` | New — supersede de AIDEC-2026-07-13-001 §5 | +| `.straymark/07-ai-audit/agent-logs/AILOG-2026-07-16-003-charter-14-durabilidad-relay.md` | New, `risk_level: medium` (cambia comportamiento del relay) | + +## Verification + +### Local checks + +```bash +cargo build --release --features test-hooks --manifest-path native/Cargo.toml +dotnet build Weft.sln -c Release + +# El contrato de IDocumentStore NO cambia: la contract suite (4 adaptadores) sigue verde sin tocar +dotnet test tests/Weft.Server.Tests -c Release +dotnet test Weft.sln -c Release + +# La carga del relay en ambos modos (evidencia del default) +dotnet run --project tests/Weft.LoadTest -c Release -- --relay + +# Convergencia real end-to-end (crítica: el reorden no debe romper clientes Yjs) +cd samples/tiptap-client && npm run check + +straymark validate --include-charters +``` + +### Production smoke (after deploy) + +No aplica: Weft es una librería. El fsync se ejercita en los tests de `FileSystemDocumentStore`; la +durabilidad ante crash de máquina real no es verificable en un shell (requiere corte de energía físico), +y NO debe clasificarse como `real_debt` si no se ejecuta aquí. + +## Risks + +- **R1 — El fsync no basta para el SLA si el store del consumidor no es durable**: probabilidad media, + severidad media. `FileSystemDocumentStore` gana fsync, pero Redis/EFCore/InMemory tienen su propia + semántica de durabilidad, ajena a Weft. + Mitigación: se documenta que la garantía de orden es *«ningún par observa lo no aceptado por el + store»*; la durabilidad física del aceptado la aporta el adaptador. El fsync cubre el store de + referencia (filesystem); los demás son responsabilidad de su backend. +- **R2 — Divergencia permanente si el fallo de append no cierra las conexiones**: probabilidad media, + severidad **alta**. Es el único modo de equivocarse: un append fallido deja el update en el doc vivo + pero nunca difundido; la próxima edición difunde solo su delta, nunca el que faltó → pares callados + para siempre. + Mitigación: `DisconnectAll` + cierre 1011 es **obligatorio** en `PersistThenBroadcast`, con un test + dedicado (`FaultyDocumentStore` que lanza en la N-ésima llamada → asertar que ningún par recibió el + delta y que las conexiones cerraron). Reconexión → SyncStep1 → el servidor reenvía. **Testeable de + forma determinista inyectando el fallo del store, sin matar procesos.** +- **R3 — Regresión de latencia de broadcast bajo un store lento**: probabilidad media, severidad media. + Persist-before-broadcast retrasa el broadcast en la latencia del append (con fsync, ~1-10ms). + Mitigación: el knob (`BroadcastThenPersist`) como válvula de escape documentada + la carga del relay + como evidencia. Sin la medición, el default sería una afirmación no medida. +- **R4 — Reorden de broadcast cross-conexión sorprende a un cliente no-Yjs**: probabilidad baja, + severidad media. + Mitigación: acotado por el pending store de yrs (updates causalmente incompletos se bufferean); el + smoke headless de `y-websocket` lo valida end-to-end; se documenta en el AIDEC y la ruta de escalada + (pump por-hub) queda registrada, no construida. +- **R5 — Un blip transitorio del store desconecta todas las conexiones de un doc**: probabilidad media, + severidad baja. + Mitigación: retry/backoff es responsabilidad del adaptador del store, no de Weft; los clientes + reconectan solos. Se documenta. +- **R6 — La doble ruta (dos modos) duplica el espacio de estados del relay**: probabilidad baja, + severidad baja. + Mitigación: ambos modos ejercitados por la misma clase de tests; el modo legacy existe como válvula + de escape, no como camino divergente (comparten todo salvo el orden de dos líneas). Candidato a borrar + `BroadcastThenPersist` en 1.0 si nadie lo usa. +- **R7 — Añadir fsync síncrono (`Flush(flushToDisk:true)` no tiene variante async en .NET) bloquea un + hilo del pool**: probabilidad media, severidad baja. + Mitigación: el append ya está fuera del turno del actor (no bloquea lecturas del doc); el coste es un + hilo del pool por append durante el fsync, que la carga del relay debe medir. Si resulta caro, evaluar + `RandomAccess.FlushToDisk` o mover a un hilo dedicado — decisión informada por la medición. + +## Tasks + +1. Sync main (con CHARTER-13), branch `charter/15-durabilidad-relay` (número de Charter 14). +2. `DurabilityMode` + `WeftServerOptions.Durability`. +3. `DocumentSession.ApplyAndCaptureDeltaAsync`. +4. `DocumentHub`: broadcast explícito por modo + fallo de append → `DisconnectAll` + relanza; wiring en + `WeftServer`. +5. `WeftConnection`: fallo del store → 1011. +6. `FileSystemDocumentStore`: fsync de archivo + directorio. +7. Cobertura de carga del relay (modo `--relay`). +8. Tests de relay (fallo inyectado, orden, reconexión, modo legacy). +9. AIDEC (supersede §5) + AILOG (`risk_level: medium`) + cerrar FU-010 + `recount`. +10. Verificación local limpia (incl. convergencia headless + carga en ambos modos) + `charter drift`. +11. Commit + push + PR. + +## Charter Closure + +1. **Atomic update (format v4)** si el drift reporta deriva no capturada. +2. **Post-merge drift check** `--range origin/main..HEAD`. +3. **Status** `in-progress` → `closed` + `closed_at`. +4. Al cerrar, `straymark followups status` debe bajar a **1 open** (FU-015). +5. **No borrar** este archivo. diff --git a/.straymark/follow-ups-backlog.md b/.straymark/follow-ups-backlog.md index 0f64996..627b54c 100644 --- a/.straymark/follow-ups-backlog.md +++ b/.straymark/follow-ups-backlog.md @@ -1,9 +1,9 @@ --- last_scan: 2026-07-15 schema_version: v1 -total_open: 2 +total_open: 1 total_promoted: 0 -total_closed_in_session: 18 +total_closed_in_session: 19 total_phase_blocked: 0 total_suspected_closed: 0 buckets: @@ -170,11 +170,11 @@ fully_extracted_ailogs: ### FU-010 — endurecimiento de durabilidad del relay: persist-before-broadcast (opcional) - **Origin**: AIDEC-2026-07-13-001 §5 (CHARTER-05) · review.md F3 (auditoría gpt-5-5 + glm-5-2) -- **Status**: open +- **Status**: closed - **Trigger**: when se requieran garantías de durabilidad duras (relay multi-nodo, o SLA de no-pérdida sin depender de la reconexión) - **Destination**: charter-replanning - **Cost**: M -- **Notes**: El relay hace **broadcast-then-persist** (AIDEC §5): aplica el update dentro del turno del actor (difunde a los pares) y persiste `IDocumentStore.AppendUpdate` **después**, fuera del turno. Un fallo del append + crash antes del snapshot pierde ese update del store; en v1 (single-node) la auto-sanación CRDT (re-sync en reconexión) lo recupera. Endurecer SOLO si se requieren semánticas duras: persist-before-broadcast (reestructurar el orden) o manejar el fallo de append cerrando la conexión + test de crash mid-operation. **Ningún gate depende hoy**; decisión consciente documentada, no un bug. +- **Notes**: El relay hace **broadcast-then-persist** (AIDEC §5): aplica el update dentro del turno del actor (difunde a los pares) y persiste `IDocumentStore.AppendUpdate` **después**, fuera del turno. Un fallo del append + crash antes del snapshot pierde ese update del store; en v1 (single-node) la auto-sanación CRDT (re-sync en reconexión) lo recupera. Endurecer SOLO si se requieren semánticas duras: persist-before-broadcast (reestructurar el orden) o manejar el fallo de append cerrando la conexión + test de crash mid-operation. **Ningún gate depende hoy**; decisión consciente documentada, no un bug. **CERRADO 2026-07-16 (CHARTER-14, AILOG-2026-07-16-003, AIDEC-2026-07-16-001 que supersede §5)**: implementado. `WeftServerOptions.Durability` con default **PersistThenBroadcast**; fallo de append → `DisconnectAll` + cierre 1011 + reconexión resincroniza; fsync de archivo+directorio en `FileSystemDocumentStore`; cobertura de carga del relay (`--relay`) que midió el coste del orden seguro ~0 en p50/p99. `BroadcastThenPersist` queda como válvula de escape opt-in. Contrato de `IDocumentStore` sin cambios. ### FU-011 — reponer la cobertura del adaptador Redis en CI (job Linux-only con service container) - **Origin**: CHARTER-06 §Scope/§Out of scope / R4 · AILOG-2026-07-13-002 §Decisions #4 (registro hand-add + recount, vía §13; el follow-up nace en tiempo de declaración de Charter, sin sección extraíble — cf. issue straymark #360) diff --git a/src/Weft.Core/Concurrency/DocumentSession.cs b/src/Weft.Core/Concurrency/DocumentSession.cs index 2cfff0b..8e815ca 100644 --- a/src/Weft.Core/Concurrency/DocumentSession.cs +++ b/src/Weft.Core/Concurrency/DocumentSession.cs @@ -89,6 +89,28 @@ public async ValueTask ApplyUpdateAsync(ReadOnlyMemory update, Cancellatio .ConfigureAwait(false); } + /// + /// Aplica un update y DEVUELVE su delta en el mismo turno del actor. Pensado para el relay con + /// persist-before-broadcast (FU-010): el delta se captura como valor de retorno —race-free frente a + /// varias conexiones concurrentes del mismo documento— para difundirlo tras persistir, en vez de + /// depender del evento (que se dispara dentro del turno, antes de + /// persistir). El delta es vacío si el update no aportó cambios nuevos (idempotente). + /// + public ValueTask ApplyAndCaptureDeltaAsync(ReadOnlyMemory update, CancellationToken ct = default) + { + ThrowIfDisposed(); + byte[] u = update.ToArray(); // copia defensiva (ver ExportUpdateSinceAsync) + return _actor.EnqueueAsync( + doc => + { + byte[] before = doc.ExportStateVector(); + doc.ApplyUpdate(u); + return doc.ExportUpdateSince(before); + }, + mutating: true, + ct); + } + /// /// Ejecuta un delegado como turno atómico respecto a las demás operaciones del mismo documento /// (transacción lógica). El recibido NO debe capturarse ni usarse fuera del diff --git a/src/Weft.Server/DocumentHub.cs b/src/Weft.Server/DocumentHub.cs index 3a768e0..a708ecc 100644 --- a/src/Weft.Server/DocumentHub.cs +++ b/src/Weft.Server/DocumentHub.cs @@ -7,22 +7,27 @@ namespace Weft.Server; /// /// Punto de encuentro de todas las conexiones de un mismo documento. Mantiene una -/// por documento (anclaje M1: el broadcast vía -/// es perezoso, se suscribe una sola vez y el refcount de sesiones -/// mantiene el documento residente mientras haya conexiones). Difunde cada update aplicado y persiste el flujo. +/// por documento (el refcount de sesiones mantiene el documento residente +/// mientras haya conexiones). Aplica cada update, lo persiste y difunde su delta; el orden entre persistir y +/// difundir lo fija (FU-010). /// internal sealed class DocumentHub : IAsyncDisposable { private readonly IDocumentStore _store; + private readonly DurabilityMode _durability; private readonly ConcurrentDictionary _connections = new(); private int _disposed; - public DocumentHub(string docId, DocumentSession session, IDocumentStore store) + public DocumentHub(string docId, DocumentSession session, IDocumentStore store, DurabilityMode durability) { DocId = docId; Session = session; _store = store; - Session.UpdateApplied += OnUpdateApplied; + _durability = durability; + // El broadcast es EXPLÍCITO (ApplyAndPersistAsync), no vía el evento UpdateApplied: capturar el + // delta como valor de retorno del turno permite ordenar append/broadcast según el modo y es + // race-free entre conexiones concurrentes del mismo documento. El evento se conserva para otros + // consumidores de DocumentSession, pero el relay ya no depende de él. } /// Identificador del documento. @@ -63,26 +68,54 @@ public void Broadcast(byte[] frame, WeftConnection? exclude) } /// - /// Aplica un update entrante al documento (turno del actor) y lo persiste. La aplicación dispara - /// , que difunde el delta a las conexiones. + /// Aplica un update entrante al documento (turno del actor), lo persiste y difunde su delta. El orden + /// entre persistir y difundir lo fija (FU-010). /// + /// + /// En , si el append falla, el update ya está en el + /// documento vivo pero NUNCA se difundió: los pares quedarían callados y desactualizados para siempre + /// (la próxima edición difunde solo su delta, no el que faltó). Por eso un fallo de append cierra + /// TODAS las conexiones del documento (el llamador las cierra con 1011): reconectan, mandan SyncStep1 + /// y el servidor —autoritativo, que sí tiene el update— reenvía el estado. Convergencia recuperada. + /// public async ValueTask ApplyAndPersistAsync(byte[] update, CancellationToken ct) { - await Session.ApplyUpdateAsync(update, ct).ConfigureAwait(false); - await _store.AppendUpdateAsync(DocId, update, ct).ConfigureAwait(false); + byte[] delta = await Session.ApplyAndCaptureDeltaAsync(update, ct).ConfigureAwait(false); + + if (_durability == DurabilityMode.BroadcastThenPersist) + { + BroadcastDelta(delta); + await _store.AppendUpdateAsync(DocId, update, ct).ConfigureAwait(false); + return; + } + + // PersistThenBroadcast (default): persistir antes de que ningún par lo vea. + try + { + await _store.AppendUpdateAsync(DocId, update, ct).ConfigureAwait(false); + } + catch + { + // El update quedó aplicado pero sin difundir: cerrar el documento entero fuerza el re-sync + // autoritativo. Relanzar deja que el llamador cierre esta conexión con 1011. + DisconnectAll(); + throw; + } + + BroadcastDelta(delta); } // El delta se difunde a TODAS las conexiones del documento. Reaplicar su propio delta en el origen es un // no-op CRDT idempotente (los clientes Yjs lo toleran), lo que evita rastrear el origen dentro del turno del // actor (que sería una carrera). El coste es un eco al emisor; aceptable para v1. - private void OnUpdateApplied(DocumentSession _, ReadOnlyMemory delta) + private void BroadcastDelta(byte[] delta) { - if (delta.IsEmpty) + if (delta.Length == 0) { return; } - Broadcast(SyncProtocol.EncodeUpdate(delta.Span), exclude: null); + Broadcast(SyncProtocol.EncodeUpdate(delta), exclude: null); } /// @@ -96,8 +129,6 @@ public async ValueTask DisposeAsync() return; } - Session.UpdateApplied -= OnUpdateApplied; - try { // Snapshot de consolidación: el estado completo dentro del turno del actor reemplaza los updates diff --git a/src/Weft.Server/Persistence/FileSystemDocumentStore.cs b/src/Weft.Server/Persistence/FileSystemDocumentStore.cs index 6136b15..9737d7d 100644 --- a/src/Weft.Server/Persistence/FileSystemDocumentStore.cs +++ b/src/Weft.Server/Persistence/FileSystemDocumentStore.cs @@ -195,10 +195,38 @@ private static async Task AtomicWriteAsync(string finalPath, ReadOnlyMemory ApplyOrCloseAsync(DocumentHub hub, byte[] payload, CancellationToken ct) + { + try + { + await hub.ApplyAndPersistAsync(payload, ct).ConfigureAwait(false); + return true; + } + catch (OperationCanceledException) + { + throw; + } + catch (WebSocketException) + { + throw; + } + catch + { + await CloseAsync(WebSocketCloseStatus.InternalServerError, "persist failed", ct) // 1011 + .ConfigureAwait(false); + return false; + } + } + private async Task DispatchAsync(DocumentHub hub, byte[] frame, CancellationToken ct) { // SyncMessage es un ref struct (span sobre el frame): se extraen tipo+payload a locales ANTES de await. @@ -142,7 +169,7 @@ private async Task DispatchAsync(DocumentHub hub, byte[] frame, Cancellati // es solo para Update en vivo, FR-019). if (Access == WeftAccess.ReadWrite) { - await hub.ApplyAndPersistAsync(payload, ct).ConfigureAwait(false); + return await ApplyOrCloseAsync(hub, payload, ct).ConfigureAwait(false); } return true; @@ -155,8 +182,7 @@ await CloseAsync(WebSocketCloseStatus.PolicyViolation, "read-only connection", c return false; } - await hub.ApplyAndPersistAsync(payload, ct).ConfigureAwait(false); // dispara UpdateApplied → broadcast - return true; + return await ApplyOrCloseAsync(hub, payload, ct).ConfigureAwait(false); case MessageType.Awareness: AwarenessProtocol.TrackClients(payload, _awarenessClients); diff --git a/src/Weft.Server/WeftServer.cs b/src/Weft.Server/WeftServer.cs index 587ef61..b1a7674 100644 --- a/src/Weft.Server/WeftServer.cs +++ b/src/Weft.Server/WeftServer.cs @@ -99,7 +99,7 @@ private async Task JoinAsync(string docId, WeftConnection connectio if (!_hubs.TryGetValue(docId, out DocumentHub? hub)) { DocumentSession session = await _broker.OpenAsync(docId, LoadDocStateAsync, ct).ConfigureAwait(false); - hub = new DocumentHub(docId, session, _store); + hub = new DocumentHub(docId, session, _store, _options.Durability); _hubs[docId] = hub; } diff --git a/src/Weft.Server/WeftServerOptions.cs b/src/Weft.Server/WeftServerOptions.cs index 40c8b44..8c4e601 100644 --- a/src/Weft.Server/WeftServerOptions.cs +++ b/src/Weft.Server/WeftServerOptions.cs @@ -30,4 +30,32 @@ public sealed class WeftServerOptions /// límite); el cliente reconecta y re-sincroniza. Por defecto 256 mensajes. /// public int MaxSendQueuePerConnection { get; set; } = 256; + + /// + /// Orden entre persistir un update y difundirlo a los pares (FU-010). Por defecto + /// : ningún par observa un update que el store no haya + /// aceptado. Ver para el trade-off con la latencia de broadcast. + /// + public DurabilityMode Durability { get; set; } = DurabilityMode.PersistThenBroadcast; +} + +/// Orden entre persistir un update en el y difundirlo (FU-010). +public enum DurabilityMode +{ + /// + /// Persistir ANTES de difundir (default). Garantiza que ningún par observa estado que el servidor no haya + /// aceptado de forma durable. Un fallo del append cierra las conexiones del documento (1011): reconectan y el + /// servidor autoritativo reenvía el estado. El coste es que el broadcast se retrasa la latencia del append (que + /// el emisor ya paga hoy en el receive loop). Nota: y-protocols no tiene ack de aplicación para Update, + /// así que el emisor nunca sabe si su update se persistió — la garantía es sobre lo que OBSERVAN los pares. + /// + PersistThenBroadcast, + + /// + /// Difundir ANTES de persistir (comportamiento heredado, CHARTER-05). Menor latencia de broadcast a costa de + /// una ventana en la que los pares tienen un update que el store aún no aceptó; en single-node la auto-sanación + /// CRDT lo recupera en la reconexión. Válvula de escape para deployments sensibles a la latencia sobre un store + /// rápido o en memoria. + /// + BroadcastThenPersist, } diff --git a/tests/Weft.LoadTest/Program.cs b/tests/Weft.LoadTest/Program.cs index 97d5474..4ea319a 100644 --- a/tests/Weft.LoadTest/Program.cs +++ b/tests/Weft.LoadTest/Program.cs @@ -1,8 +1,16 @@ using System.Collections.Concurrent; using System.Diagnostics; using Weft.Concurrency; +using Weft.LoadTest; using Weft.Yrs; +// Modo relay (FU-010/CHARTER-14): mide la latencia de broadcast del relay real en ambos modos de +// durabilidad. Distinto de la carga de US2 de abajo (que conduce el broker directamente, ciega al relay). +if (Array.IndexOf(args, "--relay") >= 0) +{ + return await RelayLoad.RunAsync(ArgInt(args, "--edits", 200)); +} + // Prueba de carga de US2/M1 (SC-006): cientos de documentos y muchas tareas concurrentes editando al // azar durante un período sostenido. Verifica (a) consistencia final de cada documento y (b) memoria // acotada — el número de documentos activos se mantiene bajo el límite pese a que el total supera el pool diff --git a/tests/Weft.LoadTest/RelayLoad.cs b/tests/Weft.LoadTest/RelayLoad.cs new file mode 100644 index 0000000..4037d0d --- /dev/null +++ b/tests/Weft.LoadTest/RelayLoad.cs @@ -0,0 +1,248 @@ +using System.Diagnostics; +using System.Net.WebSockets; +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.TestHost; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Weft; +using Weft.Server; +using Weft.Server.Auth; +using Weft.Server.Persistence; +using Weft.Server.Protocol; +using Weft.Yrs; + +namespace Weft.LoadTest; + +/// +/// Carga del relay real (FU-010/CHARTER-14): mide la latencia editor→observador de un update a través +/// del relay completo (TestServer + WebSocket + con fsync), en +/// AMBOS modos de durabilidad. Es la evidencia que respalda el default PersistThenBroadcast: sin +/// ella, la elección del default sería una afirmación no medida. La de US2 conduce +/// el broker directamente y es ciega al path de persistencia del relay. +/// +internal static class RelayLoad +{ + public static async Task RunAsync(int edits) + { + Console.WriteLine($"[relay-load] editor→observador vía relay real, {edits} ediciones/modo, FileSystemDocumentStore + fsync"); + + LatencyStats persist = await MeasureAsync(DurabilityMode.PersistThenBroadcast, edits); + LatencyStats broadcast = await MeasureAsync(DurabilityMode.BroadcastThenPersist, edits); + + Console.WriteLine($"[relay-load] PersistThenBroadcast (default): {persist}"); + Console.WriteLine($"[relay-load] BroadcastThenPersist (legacy): {broadcast}"); + Console.WriteLine( + $"[relay-load] coste del orden seguro (p50): +{persist.P50 - broadcast.P50:F1}ms, " + + $"(p99): +{persist.P99 - broadcast.P99:F1}ms"); + + // PASS = ambos modos convergieron en todas las ediciones (sin pérdidas ni timeouts). El número de + // latencia es informativo (depende del disco del runner); lo que se gatea es la corrección. + bool ok = persist.Count == edits && broadcast.Count == edits; + Console.WriteLine($"[relay-load] convergencia: persist={persist.Count}/{edits} broadcast={broadcast.Count}/{edits} " + + $"→ {(ok ? "PASS" : "FAIL")}"); + return ok ? 0 : 1; + } + + private static async Task MeasureAsync(DurabilityMode mode, int edits) + { + string dataDir = Path.Combine(Path.GetTempPath(), $"weft-relay-load-{mode}-{Environment.ProcessId}"); + Directory.CreateDirectory(dataDir); + try + { + using IHost host = await BuildHostAsync(dataDir, mode); + TestServer server = host.GetTestServer(); + + await using var editor = await RelayClient.ConnectAsync(server, "doc"); + await using var observer = await RelayClient.ConnectAsync(server, "doc"); + + var samples = new List(edits); + int converged = 0; + for (int i = 1; i <= edits; i++) + { + var sw = Stopwatch.StartNew(); + await editor.EditAsync("x"); + bool seen = await WaitUntilAsync(() => observer.TextLength >= i, TimeSpan.FromSeconds(10)); + sw.Stop(); + if (seen) + { + samples.Add(sw.Elapsed.TotalMilliseconds); + converged++; + } + } + + await host.StopAsync(); + return LatencyStats.From(samples, converged); + } + finally + { + try { Directory.Delete(dataDir, recursive: true); } catch { /* best-effort */ } + } + } + + private static async Task BuildHostAsync(string dataDir, DurabilityMode mode) + { + IHostBuilder builder = new HostBuilder().ConfigureWebHost(web => + { + web.UseTestServer(); + web.ConfigureServices(services => + { + services.AddRouting(); + services.AddWeftServer(o => + { + o.Engine = YrsEngine.Instance; + o.Durability = mode; + }); + services.AddSingleton(new AllowAll()); + services.AddSingleton(new FileSystemDocumentStore(dataDir)); + }); + web.Configure(app => + { + app.UseWebSockets(); + app.UseRouting(); + app.UseEndpoints(e => e.MapWeft("/collab")); + }); + }); + return await builder.StartAsync(); + } + + private static async Task WaitUntilAsync(Func cond, TimeSpan timeout) + { + var sw = Stopwatch.StartNew(); + while (sw.Elapsed < timeout) + { + if (cond()) + { + return true; + } + + await Task.Delay(2); + } + + return cond(); + } + + private sealed class AllowAll : IWeftAuthorizer + { + public ValueTask AuthorizeAsync(Microsoft.AspNetCore.Http.HttpContext context, string docId, CancellationToken ct) + => ValueTask.FromResult(WeftAccess.ReadWrite); + } + + private readonly record struct LatencyStats(double P50, double P99, double Max, int Count) + { + public static LatencyStats From(List samples, int count) + { + if (samples.Count == 0) + { + return new LatencyStats(0, 0, 0, count); + } + + samples.Sort(); + return new LatencyStats( + Percentile(samples, 0.50), + Percentile(samples, 0.99), + samples[^1], + count); + } + + private static double Percentile(List sorted, double p) + { + int idx = (int)Math.Ceiling(p * sorted.Count) - 1; + return sorted[Math.Clamp(idx, 0, sorted.Count - 1)]; + } + + public override string ToString() => $"p50={P50:F1}ms p99={P99:F1}ms max={Max:F1}ms (n={Count})"; + } + + /// Cliente WebSocket mínimo con un doc yrs real, hablando y-sync contra el relay. + private sealed class RelayClient : IAsyncDisposable + { + private const string Field = "body"; + private readonly WebSocket _ws; + private readonly ICrdtDoc _doc = YrsEngine.Instance.CreateDoc(); + private readonly object _lock = new(); + private readonly CancellationTokenSource _cts = new(); + private readonly Task _recv; + + private RelayClient(WebSocket ws) + { + _ws = ws; + _recv = Task.Run(ReceiveLoopAsync); + } + + public int TextLength + { + get { lock (_lock) { return _doc.GetText(Field).Length; } } + } + + public static async Task ConnectAsync(TestServer server, string docId) + { + WebSocketClient wsc = server.CreateWebSocketClient(); + WebSocket ws = await wsc.ConnectAsync(new Uri(server.BaseAddress, $"collab/{docId}"), CancellationToken.None); + var client = new RelayClient(ws); + byte[] sv; + lock (client._lock) { sv = client._doc.ExportStateVector(); } + await client.SendAsync(SyncProtocol.EncodeSyncStep1(sv)); + return client; + } + + public async Task EditAsync(string text) + { + byte[] delta; + lock (_lock) + { + byte[] before = _doc.ExportStateVector(); + _doc.InsertText(Field, 0, text); + delta = _doc.ExportUpdateSince(before); + } + + await SendAsync(SyncProtocol.EncodeUpdate(delta)); + } + + private async Task SendAsync(byte[] frame) => + await _ws.SendAsync(frame, WebSocketMessageType.Binary, endOfMessage: true, _cts.Token).ConfigureAwait(false); + + private async Task ReceiveLoopAsync() + { + var buf = new byte[16 * 1024]; + try + { + while (!_cts.IsCancellationRequested && _ws.State == WebSocketState.Open) + { + WebSocketReceiveResult r = await _ws.ReceiveAsync(buf, _cts.Token).ConfigureAwait(false); + if (r.MessageType == WebSocketMessageType.Close) + { + return; + } + + byte[] frame = buf[..r.Count]; + SyncMessage m = SyncProtocol.Decode(frame); + if (m.Type == MessageType.Sync) + { + byte[] payload = m.Payload.ToArray(); + if (m.SyncType == SyncMessageType.Step1) + { + byte[] delta; + lock (_lock) { delta = _doc.ExportUpdateSince(payload); } + await SendAsync(SyncProtocol.EncodeSyncStep2(delta)); + } + else + { + lock (_lock) { _doc.ApplyUpdate(payload); } + } + } + } + } + catch (OperationCanceledException) { } + catch (WebSocketException) { } + } + + public async ValueTask DisposeAsync() + { + await _cts.CancelAsync(); + try { await _recv; } catch { } + _ws.Dispose(); + _cts.Dispose(); + } + } +} diff --git a/tests/Weft.LoadTest/Weft.LoadTest.csproj b/tests/Weft.LoadTest/Weft.LoadTest.csproj index 3d2d094..4acecbf 100644 --- a/tests/Weft.LoadTest/Weft.LoadTest.csproj +++ b/tests/Weft.LoadTest/Weft.LoadTest.csproj @@ -10,6 +10,13 @@ + + + + + + +