feat(backend): Aenderungsfeed per Long Polling (Schritt 4)

GET /documents/{id}/changes haelt die Anfrage offen und antwortet, sobald
sich etwas tut — kumuliertes Diff seit der bekannten Version, dazu die
Ereignisse mit ihrem Absender. Ist die Basis verdichtet oder hat der Client
noch gar nichts, kommt der Volltext statt der Operationen: ein Roundtrip und
ein Sonderzustand weniger als ein eigener Fehlerpfad.

Der Feed arbeitet auf der Historie, nicht am Dokument — ein geloeschtes
Dokument muss sein DELETED noch zustellen koennen.

Blockierend auf virtuellen Threads statt DeferredResult (D76-Nachtrag 5):
So behaelt der Endpunkt die aus der Spezifikation generierte Signatur, und
API-First bleibt fuer ihn unangetastet; ein Wartender kostet trotzdem fast
nichts. Geweckt wird nach dem Commit, nie davor, und ueber einen Stempel, den
der Aufrufer VOR dem Nachsehen liest — sonst ginge ein Signal aus der Luecke
dazwischen verloren.

122 Tests. Das Szenario "ein Wartender wird geweckt" misst die Dauer: Ohne
das bestuende es auch dann, wenn der Wartende bloss in den Timeout liefe und
danach die Aenderung vorfaende. Gegenprobe: Benachrichtigung entfernt ->
genau dieses Szenario faellt; Volltext-Rueckfall entfernt -> genau jenes.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
mhoennig
2026-08-26 17:16:13 +02:00
co-authored by Claude Opus 5
parent ae85503b79
commit 741b41ef11
21 changed files with 818 additions and 25 deletions
@@ -2,6 +2,8 @@ package de.werkbaum.api
import de.werkbaum.generated.api.DocumentsApi
import de.werkbaum.generated.model.Document as ApiDocument
import de.werkbaum.generated.model.ChangeEvent as ApiChangeEvent
import de.werkbaum.generated.model.ChangeFeed as ApiChangeFeed
import de.werkbaum.generated.model.ContentPatchRequest
import de.werkbaum.generated.model.ContentPatchResult
import de.werkbaum.generated.model.DocumentCreateRequest
@@ -9,15 +11,19 @@ import de.werkbaum.generated.model.DocumentHistoryEntry as ApiHistoryEntry
import de.werkbaum.generated.model.DocumentUpdateRequest
import de.werkbaum.generated.model.RestoreRequest
import de.werkbaum.domain.ChangeAuthor
import de.werkbaum.domain.ChangeEvent
import de.werkbaum.domain.ChangeFeed
import de.werkbaum.domain.ContentPatch
import de.werkbaum.domain.Document
import de.werkbaum.domain.DocumentHistoryEntry
import de.werkbaum.service.DocumentService
import de.werkbaum.service.LiveEditingService
import org.springframework.http.CacheControl
import org.springframework.http.HttpStatus
import org.springframework.http.ResponseEntity
import org.springframework.web.bind.annotation.RequestMapping
import org.springframework.web.bind.annotation.RestController
import java.time.Duration
import java.util.UUID
/**
@@ -89,6 +95,24 @@ class DocumentsController(
)
}
/**
* Long Polling: Der Server haelt die Anfrage offen, bis sich etwas tut.
*
* `no-store` ist Pflicht ein Proxy duerfte sonst eine 204
* zwischenspeichern, und der Feed stuende still.
*/
override fun getDocumentChanges(
documentId: UUID,
since: Long,
wait: Int,
): ResponseEntity<ApiChangeFeed> {
val feed = liveEditing.changesSince(documentId, since, Duration.ofSeconds(wait.toLong()))
return ResponseEntity
.status(if (feed == null) HttpStatus.NO_CONTENT else HttpStatus.OK)
.cacheControl(CacheControl.noStore())
.body(feed?.toApi())
}
override fun getDocumentHistory(documentId: UUID): ResponseEntity<List<ApiHistoryEntry>> =
ResponseEntity.ok(service.history(documentId).map { it.toApi() })
@@ -109,6 +133,21 @@ class DocumentsController(
updatedAt = updatedAt,
)
private fun ChangeFeed.toApi(): ApiChangeFeed = ApiChangeFeed(
currentVersion = currentVersion,
events = events.map { it.toApi() },
fromVersion = fromVersion,
ops = ops?.toApi(),
content = content,
)
private fun ChangeEvent.toApi(): ApiChangeEvent = ApiChangeEvent(
version = version,
changeType = ApiChangeEvent.ChangeType.valueOf(changeType.name),
clientId = author?.clientId,
displayName = author?.displayName,
)
private fun DocumentHistoryEntry.toApi(): ApiHistoryEntry = ApiHistoryEntry(
documentId = documentId,
version = version,
@@ -0,0 +1,28 @@
package de.werkbaum.domain
import de.werkbaum.diff.LineOp
/** Was an einer Version geschehen ist und wer sie eingereicht hat. */
data class ChangeEvent(
val version: Long,
val changeType: ChangeType,
val author: ChangeAuthor? = null,
)
/**
* Alles, was seit einer bekannten Version geschehen ist.
*
* Entweder [ops] (der Normalfall der Client wendet sie an und behält Cursor
* und Scrollposition) **oder** [content] als Volltext. Letzteres, wenn die
* Basis bereits verdichtet ist oder der Client noch gar nichts hat: Dann kann
* kein exaktes Diff mehr entstehen, und über hunderte Versionen hinweg wäre
* der Cursor ohnehin nicht zu retten. Ein Roundtrip und ein Sonderzustand
* weniger als ein eigener Fehlerpfad (D76).
*/
data class ChangeFeed(
val fromVersion: Long?,
val currentVersion: Long,
val ops: List<LineOp>?,
val content: String?,
val events: List<ChangeEvent>,
)
@@ -28,6 +28,9 @@ class JpaDocumentHistoryRepository(
override fun maxVersion(documentId: UUID): Long? = jpa.maxVersion(documentId)
override fun findAfterVersion(documentId: UUID, version: Long): List<DocumentHistoryEntry> =
jpa.findByDocumentIdAndVersionGreaterThanOrderByIdAsc(documentId, version).map { it.toDomain() }
override fun findMilestones(documentId: UUID): List<DocumentHistoryEntry> =
jpa.findByDocumentIdAndMilestoneTrueOrderByIdAsc(documentId).map { it.toDomain() }
@@ -24,6 +24,11 @@ interface DocumentHistoryJpaRepository : JpaRepository<DocumentHistoryEntity, Lo
fun findByDocumentIdAndMilestoneTrueOrderByIdAsc(documentId: UUID): List<DocumentHistoryEntity>
fun findByDocumentIdAndVersionGreaterThanOrderByIdAsc(
documentId: UUID,
version: Long,
): List<DocumentHistoryEntity>
@Query("select max(e.version) from DocumentHistoryEntity e where e.documentId = :documentId")
fun maxVersion(@Param("documentId") documentId: UUID): Long?
@@ -29,6 +29,9 @@ interface DocumentHistoryRepository {
/** Höchste vergebene Versionsnummer, auch wenn deren Eintrag verdichtet wurde. */
fun maxVersion(documentId: UUID): Long?
/** Alle Einträge mit einer Version größer [version], älteste zuerst. */
fun findAfterVersion(documentId: UUID, version: Long): List<DocumentHistoryEntry>
/** Die nutzersichtbare Historie, älteste zuerst. */
fun findMilestones(documentId: UUID): List<DocumentHistoryEntry>
@@ -0,0 +1,55 @@
package de.werkbaum.service
import org.springframework.stereotype.Component
import java.time.Duration
import java.util.UUID
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.locks.ReentrantLock
import kotlin.concurrent.withLock
/**
* Weckt wartende Beobachter, wenn sich an einem Dokument etwas getan hat —
* der Kern der „Echtzeit ohne WebSocket"-Lösung (D76).
*
* **Setzt eine Einzelinstanz voraus.** Hinter einem Load Balancer erführe ein
* Beobachter auf der zweiten Instanz nichts und liefe in den Timeout. Bewusste
* Annahme, für zehn Beobachter angemessen.
*
* Gezählt wird je Dokument ein **Stempel**, nicht die Versionsnummer: Der
* Aufrufer liest ihn, **bevor** er in der Datenbank nachsieht. Ändert sich in
* der Lücke dazwischen etwas, kehrt das Warten sofort zurück, statt das Signal
* zu verpassen und die volle Wartezeit abzusitzen.
*
* Gewartet wird mit `ReentrantLock`/`Condition`, nicht mit `synchronized`
* plus `wait()`: Der Feed blockiert seinen Thread, und das soll ein
* **virtueller** Thread sein dürfen (siehe `spring.threads.virtual.enabled`) —
* ein Monitor würde dessen Träger festnageln.
*/
@Component
class ChangeNotifier {
private val lock = ReentrantLock()
private val changed = lock.newCondition()
private val stamps = ConcurrentHashMap<UUID, Long>()
fun stampOf(documentId: UUID): Long = stamps[documentId] ?: 0L
/** Meldet eine Änderung aufzurufen **nach** dem Commit, nie davor. */
fun published(documentId: UUID) = lock.withLock {
stamps.merge(documentId, 1L) { old, one -> old + one }
changed.signalAll()
}
/**
* Wartet, bis sich der Stempel des Dokuments von [since] unterscheidet.
* `true` heißt: es hat sich etwas getan. `false` heißt: Zeit abgelaufen.
*/
fun awaitChange(documentId: UUID, since: Long, timeout: Duration): Boolean = lock.withLock {
var remaining = timeout.toNanos()
while (stampOf(documentId) == since) {
if (remaining <= 0) return false
remaining = changed.awaitNanos(remaining)
}
true
}
}
@@ -8,6 +8,8 @@ import de.werkbaum.repository.DocumentHistoryRepository
import de.werkbaum.repository.DocumentRepository
import org.springframework.stereotype.Service
import org.springframework.transaction.annotation.Transactional
import org.springframework.transaction.support.TransactionSynchronization
import org.springframework.transaction.support.TransactionSynchronizationManager
import java.time.Clock
import java.time.Duration
import java.time.OffsetDateTime
@@ -20,6 +22,7 @@ class DocumentService(
private val historyRepository: DocumentHistoryRepository,
private val clock: Clock,
private val properties: LiveEditingProperties,
private val notifier: ChangeNotifier,
) {
fun findAll(): List<Document> = repository.findAll()
@@ -203,6 +206,26 @@ class DocumentService(
document.id,
document.updatedAt.minus(properties.syncRetention),
)
publishAfterCommit(document.id)
}
/**
* Weckt die Beobachter am Änderungsfeed **nach** dem Commit. Vorher
* geweckt läse ein Beobachter einen Stand, der noch nicht steht, und
* bekäme das Ereignis nie wieder. Ohne laufende Transaktion (Tests) wird
* sofort gemeldet.
*/
private fun publishAfterCommit(documentId: UUID) {
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
notifier.published(documentId)
return
}
TransactionSynchronizationManager.registerSynchronization(
object : TransactionSynchronization {
override fun afterCommit() = notifier.published(documentId)
}
)
}
/** Anlegen, Löschen, Wiederherstellen und Rückfall sind nie bloß Sync-Versionen. */
@@ -31,4 +31,12 @@ data class LiveEditingProperties(
/** Höchstlänge des Dokuments in Zeichen; der mitgelieferte Plan hat ~40 000. */
val maxContentLength: Int = 2_000_000,
/**
* Obergrenze für das Warten am Änderungsfeed. Ein Client darf keine
* beliebig lange Verbindung binden. 25 s sind gemessen: Apache auf der
* Zielumgebung hält einen Long-Poll nachweislich 30 s durch, seine
* Zeitgrenzen liegen bei 300 s (D76-Nachtrag 1/2).
*/
val maxWait: Duration = Duration.ofSeconds(25),
)
@@ -2,10 +2,14 @@ package de.werkbaum.service
import de.werkbaum.diff.DiffNotApplicableException
import de.werkbaum.diff.LineDiff
import de.werkbaum.domain.ChangeEvent
import de.werkbaum.domain.ChangeFeed
import de.werkbaum.domain.ContentPatch
import de.werkbaum.domain.ContentPatchOutcome
import de.werkbaum.domain.DocumentHistoryEntry
import de.werkbaum.repository.DocumentHistoryRepository
import org.springframework.stereotype.Service
import java.time.Duration
import java.util.UUID
import java.util.concurrent.locks.ReentrantLock
import kotlin.concurrent.withLock
@@ -32,6 +36,7 @@ class LiveEditingService(
private val documents: DocumentService,
private val history: DocumentHistoryRepository,
private val properties: LiveEditingProperties,
private val notifier: ChangeNotifier,
) {
/**
@@ -122,6 +127,66 @@ class LiveEditingService(
)
}
// -----------------------------------------------------------------------
// Änderungsfeed
// -----------------------------------------------------------------------
/**
* Alles seit [since] oder `null`, wenn innerhalb von [wait] nichts
* passiert (im Protokoll: 204).
*
* Der Feed arbeitet auf der **Historie**, nicht am Dokument: `delete()`
* entfernt das Dokument und lässt nur den Tombstone stehen ein Feed am
* Dokument müsste danach 404 liefern, ausgerechnet für das
* DELETED-Ereignis, das er zustellen soll. 404 gibt es deshalb nur bei
* gänzlich unbekannter UUID.
*/
fun changesSince(documentId: UUID, since: Long, wait: Duration): ChangeFeed? {
if (!history.exists(documentId)) throw DocumentNotFoundException(documentId)
// Den Stempel VOR dem Nachsehen lesen: Ändert sich in der Lücke
// dazwischen etwas, kehrt das Warten unten sofort zurück.
val stamp = notifier.stampOf(documentId)
feedSince(documentId, since)?.let { return it }
val timeout = minOf(wait, properties.maxWait)
if (timeout.isNegative || timeout.isZero) return null
if (!notifier.awaitChange(documentId, stamp, timeout)) return null
return feedSince(documentId, since)
}
private fun feedSince(documentId: UUID, since: Long): ChangeFeed? {
val latest = history.findLatest(documentId) ?: return null
if (latest.version <= since) return null
val events = history.findAfterVersion(documentId, since).map { it.toEvent() }
val base = history.findVersion(documentId, since)?.content
// Ohne Basis kein exaktes Diff - dann der Volltext. Das deckt den
// Nachzügler nach dem Verdichten ebenso ab wie den Erstkontakt
// (since = 0) und die PWA nach langer Offline-Zeit.
return if (base == null) {
ChangeFeed(
fromVersion = null,
currentVersion = latest.version,
ops = null,
content = latest.content,
events = events,
)
} else {
ChangeFeed(
fromVersion = since,
currentVersion = latest.version,
ops = LineDiff.compute(LineDiff.lines(base), LineDiff.lines(latest.content)),
content = null,
events = events,
)
}
}
private fun DocumentHistoryEntry.toEvent() = ChangeEvent(version, changeType, author)
private fun lockFor(documentId: UUID): ReentrantLock =
stripes[Math.floorMod(documentId.hashCode(), stripes.size)]
@@ -18,6 +18,14 @@ spring:
liquibase:
change-log: classpath:db/changelog/db.changelog-master.sql
threads:
virtual:
# Der Aenderungsfeed blockiert seinen Thread, bis sich etwas tut
# (Long Polling). Auf einem virtuellen Thread kostet das Warten fast
# nichts - ein Plattform-Thread je Beobachter waere ein Megabyte Stack
# auf einem Host, dessen knappe Groesse der Speicher ist (D76).
enabled: true
werkbaum:
live-editing:
# Schreibpause, nach der die letzte Version zum Meilenstein wird.
@@ -25,6 +33,9 @@ werkbaum:
# Danach wird eine Sync-Version verdichtet; der Feed antwortet auf ein so
# altes "since" dann mit Volltext statt mit einem Diff.
sync-retention: 1h
# Obergrenze fuers Warten am Feed. Gemessen: Apache auf der Zielumgebung
# haelt einen Long-Poll 30 s durch, seine Zeitgrenzen liegen bei 300 s.
max-wait: 25s
server:
port: 8080
+106
View File
@@ -244,6 +244,64 @@ paths:
schema:
$ref: "#/components/schemas/ProblemDetail"
/documents/{documentId}/changes:
parameters:
- name: documentId
in: path
required: true
schema:
type: string
format: uuid
get:
tags: [Documents]
operationId: getDocumentChanges
summary: Aenderungsfeed (Long Polling)
description: >
Liefert alles, was seit `since` geschehen ist. Gibt es nichts, haelt
der Server die Anfrage bis zu `wait` Sekunden offen und antwortet
sofort, sobald eine Aenderung eintrifft; sonst 204.
Der Feed arbeitet auf der **Historie**, nicht am Dokument: Ein
geloeschtes Dokument muss sein DELETED-Ereignis noch zustellen koennen.
404 gibt es deshalb nur bei gaenzlich unbekannter UUID.
Ist `since` bereits verdichtet, kann kein exaktes Diff mehr geliefert
werden - dann enthaelt die Antwort statt `ops` den **Volltext**, und
`fromVersion` fehlt. Der Client ersetzt seinen Stand.
parameters:
- name: since
in: query
required: true
description: Zuletzt gesehene Version; 0 fuer "noch nichts".
schema:
type: integer
format: int64
minimum: 0
- name: wait
in: query
required: false
description: >
Wartezeit in Sekunden. Serverseitig geklemmt - ein Client darf
keine beliebig lange Verbindung binden.
schema:
type: integer
format: int32
minimum: 0
default: 25
responses:
"200":
description: Aenderungen seit `since`
content:
application/json:
schema:
$ref: "#/components/schemas/ChangeFeed"
"204":
description: Nichts Neues innerhalb der Wartezeit
"404":
$ref: "#/components/responses/NotFound"
components:
responses:
NotFound:
@@ -448,6 +506,54 @@ components:
items:
$ref: "#/components/schemas/LineOperation"
ChangeEvent:
description: Was an einer Version geschehen ist und wer sie eingereicht hat.
type: object
required: [version, changeType]
properties:
version:
type: integer
format: int64
changeType:
type: string
enum: [CREATED, UPDATED, DELETED, RESTORED, ROLLED_BACK]
clientId:
type: string
description: Fehlt bei Aenderungen ohne Absender (etwa ueber PUT).
displayName:
type: string
description: >
Selbstgewaehlt und ohne Anmeldung eine Behauptung - in der
Oberflaeche nicht wie ein Nachweis darstellen.
ChangeFeed:
type: object
required: [currentVersion, events]
properties:
fromVersion:
type: integer
format: int64
description: >
Basis des mitgelieferten Diffs. Fehlt, wenn stattdessen `content`
geliefert wird.
currentVersion:
type: integer
format: int64
ops:
type: array
description: Kumuliertes Diff von `fromVersion` bis `currentVersion`.
items:
$ref: "#/components/schemas/LineOperation"
content:
type: string
description: >
Volltext statt Diff - wenn `since` bereits verdichtet ist oder der
Client noch gar nichts hat (`since=0`).
events:
type: array
items:
$ref: "#/components/schemas/ChangeEvent"
ProblemDetail:
type: object
description: Fehlerformat nach RFC 9457 (Problem Details)