diff --git a/agent-client/src/main.rs b/agent-client/src/main.rs index 9d07e803..cf89f879 100644 --- a/agent-client/src/main.rs +++ b/agent-client/src/main.rs @@ -616,6 +616,7 @@ pub fn msg_name(msg: &onlinerpg_shared::ServerMessage) -> &'static str { ServerMessage::EstateChestState { .. } => "EstateChestState", ServerMessage::FurniturePurchaseResult { .. } => "FurniturePurchaseResult", ServerMessage::FurnitureSelectionNotice { .. } => "FurnitureSelectionNotice", + ServerMessage::Pong { .. } => "Pong", } } diff --git a/client/src/lib/components/PlayerControl.svelte b/client/src/lib/components/PlayerControl.svelte index 72e1d4a0..6876bd2c 100644 --- a/client/src/lib/components/PlayerControl.svelte +++ b/client/src/lib/components/PlayerControl.svelte @@ -244,7 +244,8 @@ () => networkManager.nextMoveRequestId(), (goal) => networkManager.sendMoveGoal(goal), (requestId) => networkManager.sendMoveStop(requestId), - (input) => networkManager.sendMoveDirection(input) + (input) => networkManager.sendMoveDirection(input), + () => networkManager.measuredOneWayMs() ) const { renderer } = useThrelte() diff --git a/client/src/lib/components/player-control/server-movement.test.ts b/client/src/lib/components/player-control/server-movement.test.ts index 7d07155d..6aefb9dd 100644 --- a/client/src/lib/components/player-control/server-movement.test.ts +++ b/client/src/lib/components/player-control/server-movement.test.ts @@ -7,12 +7,18 @@ import { buildAttackState } from './player-state-builders' import { transitionAttackToIdle } from './fsm/combat' import type { PlayerState } from '../../utils/movementUtils' -function setup() { +function setup(oneWayMs = 0) { let id = 0 const goal = vi.fn() const stop = vi.fn() const direction = vi.fn() - const movement = new ServerMovement(() => ++id, goal, stop, direction) + const movement = new ServerMovement( + () => ++id, + goal, + stop, + direction, + () => oneWayMs + ) return { movement, goal, stop, direction } } @@ -138,12 +144,31 @@ describe('keyboard release', () => { expect(direction.mock.lastCall?.[0].request_id).toBe(3) expect(movement.stopping).toBe(false) expect(movement.acceptStopped(stopped(2, 150))).toBe(false) + // The restart never stands still when the server still holds a route: + // with the stop unacknowledged that route is the truth, and it must not + // run past what was approved. With the stop acknowledged the server has + // already halted the avatar, so there is nothing left to retain. vi.advanceTimersByTime(200) - expect(movement.sample(() => false)).toBeNull() + const pose = movement.sample(() => false) + if (acknowledged) { + expect(pose).toBeNull() + } else { + expect(pose!.position.x).toBeLessThanOrEqual(1.5) + expect(pose!.position.x).toBeGreaterThan(0) + } expect( movement.acceptPath({ ...directionPath(200), request_id: 3 }) ).toBe(true) - expect(movement.sample(() => false)?.position.x).toBeCloseTo(0.6, 5) + if (acknowledged) { + // Nothing was retained, so the new approval opens a fresh clock and + // the pose is the one the server stamped. + expect(movement.sample(() => false)?.position.x).toBeCloseTo(0.6, 5) + } else { + // The replacement joins the retained route's playback clock, so the + // avatar carries on from where it already was instead of snapping + // back to the instant the server stamped. + expect(movement.sample(() => false)?.position.x).toBeCloseTo(0.9, 5) + } } ) @@ -553,7 +578,7 @@ describe('server approved movement', () => { expect(movement.sample(() => false)?.position.x).toBeCloseTo(0.45) }) - it.each(['stop', 'direction', 'finish', 'relocate'])( + it.each(['stop', 'finish', 'relocate'])( 'discards queued drag updates and in-flight replies after %s', (action) => { const { movement, goal } = setup() @@ -565,14 +590,7 @@ describe('server approved movement', () => { movement.request(9, 3, false) if (action === 'stop') movement.clear() else if (action === 'relocate') movement.clear(false) - else if (action === 'finish') movement.finish() - else - movement.direction({ - rotation: 0, - forward: 1, - turn: 0, - sprinting: false, - }) + else movement.finish() vi.advanceTimersByTime(200) expect(movement.sample(() => false)).toBeNull() expect(goal).toHaveBeenCalledTimes(2) @@ -585,6 +603,24 @@ describe('server approved movement', () => { } ) + it('keeps the approved drag route but rejects its replies after direction', () => { + const { movement } = setup() + movement.request(3, 3, false) + vi.advanceTimersByTime(250) + movement.acceptPath(path()) + movement.request(6, 3, false) + movement.acceptPath({ ...path(2), server_time_ms: 100 }) + movement.request(9, 3, false) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + // The server still holds the approved route until the new input lands, so + // playback continues rather than stalling. What must not come back is a + // stale reply resurrecting the path the player just replaced. + vi.advanceTimersByTime(200) + expect(movement.sample(() => false)).not.toBeNull() + expect(movement.acceptPath({ ...path(2), server_time_ms: 200 })).toBe(false) + expect(movement.acceptPath({ ...path(3), server_time_ms: 300 })).toBe(false) + }) + it.each([20, 75, 125])( 'keeps keyboard motion continuous with %ims initial latency and jitter', (latency) => { @@ -703,7 +739,7 @@ describe('server approved movement', () => { expect(movement.sample(() => false)?.position.x).toBeCloseTo(1.5, 5) }) - it.each(['stop', 'click', 'direction', 'finish'])( + it.each(['stop', 'click', 'finish'])( 'discards queued keyboard paths after %s', (action) => { const { movement } = setup() @@ -714,20 +750,28 @@ describe('server approved movement', () => { movement.acceptPath(directionPath(100)) if (action === 'stop') movement.clear() else if (action === 'click') movement.request(3, 0, false) - else if (action === 'finish') movement.finish() - else - movement.direction({ - rotation: 1, - forward: 1, - turn: 0, - sprinting: false, - }) + else movement.finish() vi.advanceTimersByTime(150) expect(movement.sample(() => false)).toBeNull() expect(movement.acceptPath(directionPath(200))).toBe(false) } ) + it('keeps playing the approved route after a direction change', () => { + const { movement } = setup() + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + vi.advanceTimersByTime(250) + movement.acceptPath(directionPath(0)) + vi.advanceTimersByTime(50) + movement.acceptPath(directionPath(100)) + movement.direction({ rotation: 1, forward: 1, turn: 0, sprinting: false }) + vi.advanceTimersByTime(150) + // The server has not seen the new heading yet, so the route it still + // holds is the truthful position to show — and its replies stay rejected. + expect(movement.sample(() => false)).not.toBeNull() + expect(movement.acceptPath(directionPath(200))).toBe(false) + }) + it('rejects samples older than a queued keyboard path', () => { const { movement } = setup() movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) @@ -1025,3 +1069,147 @@ it('displays approved turn timing even when position does not change', () => { expect(pose?.position).toEqual(approved.position) expect(pose?.rotation).toBeCloseTo(Math.PI / 4) }) + +describe('link compensation', () => { + // One coherent timeline: the client and server clocks share an origin, the + // request leaves at local 0, the server stamps the reply at local rtt/2, and + // the reply lands at local rtt. The avatar should then sit exactly where + // the server is at any later local time — not one round trip behind it. + function replyInFlight(rtt: number) { + const { movement } = setup(rtt / 2) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + vi.advanceTimersByTime(rtt) + movement.acceptPath(directionPath(rtt / 2)) + return movement + } + + it.each([20, 150, 250])( + 'renders the server present, not its past, at %ims RTT', + (rtt) => { + const movement = replyInFlight(rtt) + expect(movement.sample(() => false)?.position.x).toBeCloseTo( + rtt * 0.003, + 4 + ) + vi.advanceTimersByTime(100) + expect(movement.sample(() => false)?.position.x).toBeCloseTo( + (rtt + 100) * 0.003, + 4 + ) + } + ) + + it('closes exactly the one-way gap the shift used to leave', () => { + // Both clients see the same reply at the same instant. The one that + // measured its link shows where the server is; the one that could not + // shows where the server was one way trip ago — at 3 m/s, 0.375 m over + // a 125ms half of the 250ms round trip. + const { movement: measured } = setup(125) + const { movement: blind } = setup(0) + const input = { rotation: 0, forward: 1, turn: 0, sprinting: false } + measured.direction(input) + blind.direction(input) + vi.advanceTimersByTime(250) + measured.acceptPath(directionPath(125)) + blind.acceptPath(directionPath(125)) + + const shown = measured.sample(() => false)!.position.x + const stale = blind.sample(() => false)!.position.x + expect(shown).toBeCloseTo(0.75, 4) + expect(stale).toBeCloseTo(0.375, 4) + expect(shown - stale).toBeCloseTo(0.375, 4) + }) + + it('leaves a link that measured nothing anchored on arrival', () => { + const { movement } = setup(0) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + vi.advanceTimersByTime(20) + movement.acceptPath(directionPath(20)) + expect(movement.sample(() => false)?.position.x).toBeCloseTo(0.06, 4) + }) + + it('refuses to shift further than the approved path reaches', () => { + // A pathological probe must not extrapolate past the approved route. + const { movement } = setup(5000) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + const approved = directionPath(1000) + approved.waypoints = [ + { + position: { x: 4, y: 0, z: 0 }, + floor_level: 0, + rotation: Math.PI / 2, + travel_seconds: 1, + }, + ] + movement.acceptPath(approved) + vi.advanceTimersByTime(10_000) + const pose = movement.sample(() => false) + expect(pose?.position.x).toBeLessThanOrEqual(4) + expect(pose?.speed).toBe(0) + }) + + // Before this, a direction change nulled the anchor, so the character stood + // motionless for a full round trip. The retained route is not a prediction: + // the server still holds that route until the new input lands, so showing + // it is simply showing the truth a moment longer. + it.each([150, 250, 400])( + 'keeps moving through an input change at %ims RTT', + (rtt) => { + const { movement, direction } = setup(rtt / 2) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + movement.acceptPath(directionPath(1000)) + vi.advanceTimersByTime(100) + const before = movement.sample(() => false) + expect(before?.position.x).toBeGreaterThan(0) + + // Turn: a new request leaves, but the replacement has not come back. + movement.direction({ rotation: 0, forward: 0, turn: 1, sprinting: false }) + expect(direction).toHaveBeenCalledTimes(2) + vi.advanceTimersByTime(rtt / 2) + const during = movement.sample(() => false) + expect(during).not.toBeNull() + expect(during?.position.x).toBeGreaterThanOrEqual(before!.position.x) + expect(movement.active).toBe(true) + } + ) + + it('still starts from a standstill when no path is approved yet', () => { + const { movement } = setup(125) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + expect(movement.sample(() => false)).toBeNull() + movement.direction({ rotation: 0, forward: 1, turn: -1, sprinting: false }) + expect(movement.sample(() => false)).toBeNull() + }) + + it('drops the retained route when the caller stops outright', () => { + const { movement } = setup(125) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + movement.acceptPath(directionPath(1000)) + vi.advanceTimersByTime(100) + movement.clear() + expect(movement.sample(() => false)).toBeNull() + }) + + it('re-anchors on the replacement path once it is approved', () => { + const { movement } = setup(125) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + movement.acceptPath(directionPath(1000)) + vi.advanceTimersByTime(100) + movement.direction({ rotation: 0, forward: 0, turn: 1, sprinting: false }) + expect(movement.sample(() => false)).not.toBeNull() + // The approval for the new request takes over from the retained route. + expect(movement.acceptPath({ ...directionPath(1350), request_id: 2 })).toBe( + true + ) + expect(movement.isCurrentRequest(2)).toBe(true) + }) + + it('ignores a replacement path for the request it already replaced', () => { + const { movement } = setup(125) + movement.direction({ rotation: 0, forward: 1, turn: 0, sprinting: false }) + movement.acceptPath(directionPath(1000)) + movement.direction({ rotation: 0, forward: 0, turn: 1, sprinting: false }) + // An approval for the pre-turn request must not resurrect the old route. + expect(movement.acceptPath(directionPath(1200))).toBe(false) + }) +}) diff --git a/client/src/lib/components/player-control/server-movement.ts b/client/src/lib/components/player-control/server-movement.ts index 41c8c4c6..64af9033 100644 --- a/client/src/lib/components/player-control/server-movement.ts +++ b/client/src/lib/components/player-control/server-movement.ts @@ -11,6 +11,9 @@ import { shortestWrappedDeltaX, wrapWorldX } from '../../terrain/world-wrap' const SEND_INTERVAL_MS = 200 const MAX_EXTRAPOLATION_MS = 500 const STOP_BLEND_MS = 120 +/// Cap on the anchor shift. A probe that reports a pathological round trip +/// should leave the character slow and trailing, never teleported. +const MAX_SHIFT_MS = 250 type Pose = { position: Position @@ -44,7 +47,8 @@ export class ServerMovement { private readonly nextId: () => number, private readonly sendGoal: (goal: MoveGoal) => void, private readonly sendStop: (requestId: number) => void, - private readonly sendDirection: (input: MoveDirection) => void = () => {} + private readonly sendDirection: (input: MoveDirection) => void = () => {}, + private readonly oneWayMs: () => number = () => 0 ) {} get active() { @@ -95,13 +99,37 @@ export class ServerMovement { previous.turn !== input.turn || previous.sprinting !== input.sprinting ) { - this.clear(false) + this.linger() this.requestId = this.nextId() this.directionInput = { ...input, request_id: this.requestId } } this.sendDirection(this.directionInput!) } + /** + * Input changed under a live approved path. Everything about the request + * is stale, but the path itself is not: keeping it playing means a distant + * client sees the new heading arrive late instead of stopping dead for a + * full round trip. Only a client with no approved path yet starts still. + */ + private linger() { + if (this.timer !== null) clearTimeout(this.timer) + this.timer = null + this.pending = null + this.stopBlend = null + this.requestId = null + this.directionInput = null + this.sentGoalIds = [] + this.blockedPose = null + // A fresh direction supersedes a stop still in flight, so its reply must + // not be allowed to blend the avatar to a halt mid-stride. + this.stopId = null + if (!this.anchor) { + this.waypoints = [] + this.queuedUpdates = [] + } + } + stopDirection() { if (!this.directionInput || this.stopping) return this.stopId = this.nextId() @@ -136,10 +164,15 @@ export class ServerMovement { return false if (goalIndex >= 0) this.sentGoalIds.splice(0, goalIndex) const now = performance.now() - // Retargets and progress share the approved path's playback clock. + // Retargets and progress share the approved path's playback clock. The + // first anchor is dated back by the measured one-way delay: the pose in + // the message describes where the player was when the server stamped it, + // and anchoring on arrival would trail the authoritative position by a + // full one-way trip for as long as the walk lasts. Later updates inherit + // the same offset through their server-time delta. const at = this.anchor ? this.anchorAt + (progress.server_time_ms - this.anchor.server_time_ms) - : now + : now - this.shiftMs() const update = { progress, waypoints, at } if (at > now) this.queuedUpdates.push(update) else { @@ -149,6 +182,13 @@ export class ServerMovement { return true } + private shiftMs(): number { + const measured = this.oneWayMs() + return Number.isFinite(measured) + ? Math.min(Math.max(measured, 0), MAX_SHIFT_MS) + : 0 + } + private get latestServerTime() { return ( this.queuedUpdates.at(-1)?.progress.server_time_ms ?? diff --git a/client/src/lib/network/linkLatency.test.ts b/client/src/lib/network/linkLatency.test.ts new file mode 100644 index 00000000..ed173f0c --- /dev/null +++ b/client/src/lib/network/linkLatency.test.ts @@ -0,0 +1,156 @@ +import { describe, expect, it } from 'vitest' +import { LinkLatency } from './linkLatency' + +/** Drives a tracker through a link with a fixed one-way cost per direction. */ +function harness(options: { rttMs: number }) { + const link = new LinkLatency() + let rttMs = options.rttMs + let seq = 0 + let now = 10_000 + const sent: { seq: number; at: number }[] = [] + + const pump = () => { + link.due(now, (s, clientTimeMs) => { + seq = s + sent.push({ seq: s, at: now }) + expect(clientTimeMs).toBe(Math.round(now)) + }) + } + + const answer = (atSeq = seq) => { + const sentAt = sent.find((p) => p.seq === atSeq)?.at + if (sentAt === undefined) return false + return link.accept(atSeq, Math.round(sentAt), sentAt + rttMs) + } + + return { + link, + pump, + answer, + setRtt: (ms: number) => { + rttMs = ms + }, + advance: (ms: number) => { + now += ms + }, + now: () => now, + } +} + +describe('LinkLatency', () => { + it('reports nothing until a probe is answered', () => { + const { link, pump } = harness({ rttMs: 200 }) + pump() + expect(link.rttMs).toBe(0) + expect(link.oneWayMs).toBe(0) + expect(link.sample).toBeNull() + }) + + it.each([20, 150, 250, 400])('halves a %ims round trip', (rttMs) => { + const h = harness({ rttMs }) + h.pump() + expect(h.answer()).toBe(true) + expect(h.link.rttMs).toBeCloseTo(rttMs, 5) + expect(h.link.oneWayMs).toBeCloseTo(rttMs / 2, 5) + }) + + it('paces probes so one player is not a message flood', () => { + const h = harness({ rttMs: 100 }) + h.pump() + h.answer() + for (let i = 0; i < 50; i++) { + h.advance(10) + h.pump() + } + expect(h.link.rttMs).toBeCloseTo(100, 5) + }) + + it('rejects an answer whose echoed time does not match the probe', () => { + const h = harness({ rttMs: 100 }) + h.pump() + expect(h.link.accept(999, 10_000, 10_100)).toBe(false) + expect(h.link.rttMs).toBe(0) + }) + + it('rejects a reordered answer for a probe it already retired', () => { + const h = harness({ rttMs: 100 }) + h.pump() + h.answer() + const stale = h.link.rttMs + expect(h.link.accept(1, 10_000, 10_400)).toBe(false) + expect(h.link.rttMs).toBe(stale) + }) + + it('folds a slow sample in gradually rather than jumping the estimate', () => { + const h = harness({ rttMs: 100 }) + h.pump() + h.answer() + expect(h.link.rttMs).toBeCloseTo(100, 5) + // A route change to a much worse path: the estimate must move, but only + // part of the way, so one bad sample cannot inflate playback lag. + h.setRtt(400) + h.advance(2100) + h.pump() + h.answer() + expect(h.link.rttMs).toBeGreaterThan(100) + expect(h.link.rttMs).toBeLessThan(400) + expect(h.link.rttMs).toBeCloseTo(175, 5) + }) + + it('recovers from a spike instead of staying inflated', () => { + const link = new LinkLatency() + let now = 0 + const sentAt: number[] = [] + const pump = () => + link.due(now, () => { + sentAt.push(now) + }) + const answer = (rtt: number) => { + const at = sentAt.at(-1)! + link.accept(sentAt.length, at, at + rtt) + } + for (const rtt of [100, 100, 100, 900, 100, 100, 100, 100, 100, 100]) { + pump() + answer(rtt) + now += 2100 + } + // Smoothing with alpha 0.25 leaves the estimate well below the spike. + expect(link.rttMs).toBeLessThan(300) + expect(link.jitterMs).toBeGreaterThan(0) + }) + + it('forgets everything on reset, as a new socket requires', () => { + const h = harness({ rttMs: 250 }) + h.pump() + h.answer() + expect(h.link.rttMs).toBeGreaterThan(0) + h.link.reset() + expect(h.link.rttMs).toBe(0) + expect(h.link.sample).toBeNull() + }) + + it('lets a probe go out again immediately after a reset', () => { + const h = harness({ rttMs: 250 }) + h.pump() + h.link.reset() + h.pump() + expect(h.answer()).toBe(true) + }) + + it('reports spread across the window as jitter', () => { + const link = new LinkLatency() + const sentAt: number[] = [] + let now = 0 + const pump = () => link.due(now, () => void sentAt.push(now)) + const answer = (rtt: number) => { + const at = sentAt.at(-1)! + link.accept(sentAt.length, at, at + rtt) + } + for (const rtt of [100, 100, 100]) { + pump() + answer(rtt) + now += 2100 + } + expect(link.jitterMs).toBe(0) + }) +}) diff --git a/client/src/lib/network/linkLatency.ts b/client/src/lib/network/linkLatency.ts new file mode 100644 index 00000000..264ab002 --- /dev/null +++ b/client/src/lib/network/linkLatency.ts @@ -0,0 +1,121 @@ +/// How often a probe goes out. Fast enough to notice a route change (a roaming +/// client, a VPN toggle), rare enough that a 5,000-player server sees one +/// message per player per interval rather than a flood. +const PROBE_INTERVAL_MS = 2000 + +/// A probe older than this never produced an answer; count it as loss so a +/// link that only answers half the time does not look idle-fast. +const PROBE_TIMEOUT_MS = 6000 + +/// Samples kept for the floor. The minimum round trip over a window is the +/// least-delayed sample the link produced, which is the honest basis for +/// compensation: a spike should widen the error band, not move the estimate. +const WINDOW = 8 + +/// Time constant for the smoothed estimate. Fast enough to follow a route +/// change within a couple of probes, slow enough that one spike does not +/// drag the whole estimate up. +const SMOOTHING = 0.25 + +export type LinkSample = { + rttMs: number + /** Spread between the fastest and slowest sample in the window. */ + jitterMs: number +} + +/** + * Round-trip estimator for the websocket link. + * + * The client cannot otherwise know how much of a delay is spent in each + * direction, and that is what playback needs: `oneWayMs` is the amount the + * anchor is dated back by. Anchoring on arrival instead trails the + * authoritative position by a whole one-way trip for as long as the walk + * lasts, which is the difference between smooth and rubber-banding on a link + * that crosses an ocean. + * + * Deliberately not tracking a clock offset. `server_time_ms` is an epoch + * stamp while the local clock is monotonic from an arbitrary origin, so the + * difference between them is a large constant rather than a usable offset; + * measuring against a live server confirmed it. Until something needs a real + * clock sync, the round trip is the honest and sufficient number. + */ +export class LinkLatency { + private rtts: number[] = [] + private smoothedRtt = 0 + private seeded = false + private nextProbeAt = 0 + private seq = 0 + private readonly pending = new Map() + + /** Round-trip estimate, or 0 until the first answer lands. */ + get rttMs(): number { + return this.seeded ? this.smoothedRtt : 0 + } + + /** Half the round trip: the best available estimate of one direction. */ + get oneWayMs(): number { + return this.rttMs / 2 + } + + /** Spread of the recent window. A wide spread means corrections will show. */ + get jitterMs(): number { + if (this.rtts.length < 2) return 0 + return Math.max(...this.rtts) - Math.min(...this.rtts) + } + + get sample(): LinkSample | null { + return this.seeded + ? { rttMs: this.smoothedRtt, jitterMs: this.jitterMs } + : null + } + + /** + * A probe is due. `now` is the local monotonic clock; the caller supplies + * the send so a probe is only claimed when it really goes on the wire. + */ + due(now: number, send: (seq: number, clientTimeMs: number) => void): boolean { + this.expire(now) + if (now < this.nextProbeAt) return false + this.nextProbeAt = now + PROBE_INTERVAL_MS + const seq = ++this.seq + this.pending.set(seq, now) + send(seq, Math.round(now)) + return true + } + + /** + * Fold an answer in. `clientTimeMs` is echoed back so a reordered or stale + * answer cannot be mistaken for the probe it claims to be; an unknown seq + * is dropped rather than guessed at. + */ + accept(seq: number, clientTimeMs: number, now: number): boolean { + const sentAt = this.pending.get(seq) + if (sentAt === undefined || Math.round(sentAt) !== clientTimeMs) + return false + this.pending.delete(seq) + const rtt = now - sentAt + if (rtt < 0) return false + if (this.seeded) this.smoothedRtt += (rtt - this.smoothedRtt) * SMOOTHING + else { + this.smoothedRtt = rtt + this.seeded = true + } + this.rtts.push(rtt) + if (this.rtts.length > WINDOW) this.rtts.shift() + return true + } + + reset() { + this.rtts.length = 0 + this.pending.clear() + this.smoothedRtt = 0 + this.seeded = false + this.nextProbeAt = 0 + } + + private expire(now: number) { + for (const [seq, sentAt] of this.pending) { + if (now - sentAt >= PROBE_TIMEOUT_MS) this.pending.delete(seq) + } + } +} diff --git a/client/src/lib/network/networkTypes.ts b/client/src/lib/network/networkTypes.ts index 4e5cb926..077eb18a 100644 --- a/client/src/lib/network/networkTypes.ts +++ b/client/src/lib/network/networkTypes.ts @@ -251,6 +251,7 @@ export type ClientMessage = | { InteractObject: { object_type: string; object_id: number } } | 'StopInteraction' | 'Heartbeat' + | { Ping: { seq: number; client_time_ms: number } } | 'ResyncWorld' | { EquipItem: { instance_id: number } } | { SelectAmmo: { item_def_id: string | null } } diff --git a/client/src/lib/network/socket.ts b/client/src/lib/network/socket.ts index b4b66ed7..390ae0f5 100644 --- a/client/src/lib/network/socket.ts +++ b/client/src/lib/network/socket.ts @@ -55,6 +55,7 @@ import initWasm, { import { createEvent } from './networkEvents' import { handleServerMessage, resetTerrainDownloads } from './messageHandlers' import { worldView } from './worldView' +import { LinkLatency } from './linkLatency' import type { AccountCharacter, CharacterClass, @@ -124,6 +125,10 @@ class NetworkManager { private lastServerUrl: string = '' private lastCharacterId: number | null = null private wasmReady = false + /// Measured link cost. Movement playback dates its anchor off this, so a + /// player on another continent renders the server's present rather than its + /// past. Per connection: a new socket is a new route. + readonly link = new LinkLatency() /// Reset per socket: the handshake is per connection, not per session. private handshakeSent = false /// Server refused this build at the handshake. Reconnecting cannot fix a @@ -243,6 +248,7 @@ class NetworkManager { console.log('Attempting to connect to:', targetUrl) this.handshakeSent = false this.lastAuthErrorMessage = null + this.link.reset() this.socket = new WebSocket(targetUrl) this.socket.binaryType = 'arraybuffer' @@ -252,7 +258,10 @@ class NetworkManager { // send path calls ensureHandshake() too, so a first message that races // this callback still goes out second. void this.ensureWasm().then(() => { - if (this.socket?.readyState === WebSocket.OPEN) this.ensureHandshake() + if (this.socket?.readyState === WebSocket.OPEN) { + this.ensureHandshake() + this.pumpProbe() + } }) gameStore.update((state) => ({ ...state, isConnected: true })) serverNotice.set(null) @@ -307,6 +316,11 @@ class NetworkManager { try { const bytes = new Uint8Array(event.data as ArrayBuffer) const message = deserialize_server_message(bytes) + // Traffic is the clock a busy client runs on, so every inbound frame + // is a chance to keep the probe going. Without this an idle player + // would only ever measure the one link sample taken at connect. + this.pumpProbe() + if (this.acceptPong(message)) return handleServerMessage( message, this.messageEvents, @@ -327,6 +341,24 @@ class NetworkManager { } } + /// Keeps the round-trip estimate current. Cheap: one message per interval, + /// and only while the socket is actually open. + private pumpProbe() { + this.link.due(performance.now(), (seq, clientTimeMs) => { + this.sendMessage({ Ping: { seq, client_time_ms: clientTimeMs } }) + }) + } + + private acceptPong(message: unknown): boolean { + if (!message || typeof message !== 'object' || !('Pong' in message)) + return false + const { seq, client_time_ms } = ( + message as { Pong: { seq: number; client_time_ms: number } } + ).Pong + this.link.accept(seq, client_time_ms, performance.now()) + return true + } + /// Full jitter: half the capped delay plus a random half, so clients that /// dropped together don't come back together. private reconnectDelay(): number { @@ -404,6 +436,11 @@ class NetworkManager { return this.socket?.readyState === WebSocket.OPEN && this.wasmReady } + /// Half the measured round trip, in ms. 0 until the first probe answers. + measuredOneWayMs(): number { + return this.link.oneWayMs + } + private sendAndSerialize(msg: ClientMessage): boolean { if (!this.isConnected()) return false this.ensureHandshake() diff --git a/doc/SERVER_PLAYER_MOVEMENT.md b/doc/SERVER_PLAYER_MOVEMENT.md index bb822e32..5ca95491 100644 --- a/doc/SERVER_PLAYER_MOVEMENT.md +++ b/doc/SERVER_PLAYER_MOVEMENT.md @@ -1,5 +1,9 @@ # 서버 주도 플레이어 이동: 설계와 검증 기록 +2026-09-26 갱신. 프로토콜 104로 **링크 지연 보정**을 추가했다(왕복 탐색, 편도 +지연만큼의 앵커 보정, 입력 변경 시 승인 구간 유지). 아래 "링크 지연 보정" 참조. +로컬 경로 예측은 여전히 없다. + 2026-09-25 구현 갱신. **서버 승인 후 이동하는 1단계로 전체 플레이어 이동을 통합했다.** 브라우저와 agent-client는 목표·방향·정지를 요청하고 서버가 위치·높이·층·충돌을 결정한다. 예측 출발은 구현하지 않았다. 아래 설계 논의와 기존 A* 벤치마크는 비교 근거로 보존한다. @@ -640,6 +644,76 @@ sequenceDiagram 연타 중 로컬 A*도 최신 목표로 합치고 오래된 결과를 무효화한다. 예측 이동의 충돌 검사는 필요한 구간을 나눠 처리하며, 승인 전후의 진행 차이와 프레임 정지를 함께 검증한다. +## 링크 지연 보정 (프로토콜 104) + +2026-09-26. 위 2단계(예측 출발) 도입 전제로, **먼 서버에 대한 표시 지연만 제거하는 +기초 계층**을 추가했다. 로컬 경로 예측은 여전히 없다. + +### 무엇을 측정했는가 + +1단계까지 클라는 자기 왕복 시간을 알지 못했다. `Heartbeat`은 무응답짜리 연결 +유지 신호였고, 이동 응답의 `server_time_ms`도 클라 시계와의 오프셋을 알려주지 않는다. +따라서 표시가 얼마나 뒤처졌는지 **추정할 기준 자체가 없었다.** + +- 프로토콜 **104**: `ClientMessage::Ping { seq, client_time_ms }`와 + `ServerMessage::Pong { seq, client_time_ms, server_time_ms }`를 추가했다. + `Pong`는 `ClientMessage::Ping`보다 앞선 대기로 처리되며, 클라는 두 시계 중 어느 것도 + 믿지 않고 왕복을 잰다. 두 열거형은 MessagePack가 **정수 인덱스**로 부호화하므로 + 새 변수는 반드시 각 열거형의 **끝**에 둔다. +- 클라 `linkLatency.ts`는 2초마다 1회, 응답이 오면 RTT를 EWMA(α=0.25)로 접합하고 + 최근 8개 표본의 최소값과 최대값 차이를 지터로 보고한다. +- **시계 오프셋은 의도적으로 추적하지 않는다.** `Pong`의 `server_time_ms`는 epoch + 시각이고 클라 `client_time_ms`는 임의 원점의 단조 시계라, 두 값의 차이는 쓸 수 있는 + 오프셋이 아니라 큰 상수다. 실제 서버에 붙여 측정해 확인했다(§검증). 왕복 시간만이 + 보정에 필요한 숫자이므로 오프셋 machinery를 싣지 않는다. `server_time_ms`는 이후 + 진짜 시계 동기화가 필요해질 때를 위해 프로토콜에 남겨둔다. +- 클라는 후속 `PlayerMovePath`/`PlayerMoveProgress`를 받을 때마다 다음 탐색를 예약한다. + 유휴 상태에서도 2초 간격은 유지된다. 5,000인 서버 기준 인당 2초당 1메시지다. + +### 무엇이 바뀌었는가 + +| 문제 | 1단계 동작 | 104 이후 | +|---|---|---| +| 입력 변경 시 정지 | 방향 입력이 바뀌면 `clear(false)`가 앵커를 `null`로 만들어 `sample()`이 `null`을 돌려줬다. 원격 플레이어는 RTT만큼 **멈춰 서 있었다** | 승인된 구간을 계속 재생한다. 새 입력 전까지 서버가 가진 그대로가 **참**이므로 예측이 아니라 사실 표시다. 후속 응답은 여전히 거부된다 | +| 한 방향 지연 | 앵커가 **도착 시각**에 걸려, 재생 위치가 항상 편도 지연만큼 뒤처졌다 | 첫 앵커를 측정된 편도 지연만큼 **뒤로** dating한다. 이후 갱신은 `server_time_ms` 차이로 같은 오프셋을 이어받는다 | +| 표정 시계 | 500ms 고정 | 편도 지연 최대 250ms로 제한. 탐색가 경로 앞을 넘어도 승인 구간에서 멈춘다 | + +### 효과가 없는 경우와 남은 것 + +- **탐색 전까지는 지연 0이다.** 게임을 시작하자마자 2초는 1단계와 같다. 첫 앵커가 + 생기면 그때부터 보정된다. +- **입력→첫 모션은 여전히 1 RTT다.** 250ms RTT에서 250ms 지연은 남아 있다. 이것을 + 없애는 것은 2단계 예측 출발이며, 아래 설계대로 별도 과제로 다룬다. +- **로컬 A\*과 충돌 예측은 없다.** 보정은 표시 시계만 건드린다. +- **원격 플레이어는 보간하지 않는다.** `remotePlayerManager`는 200ms 갱신 + 편도 + 지연으로 늦고 5Hz 계단 형태다. 자기 이동이 아니라 상대 표시 문제다. +- **전투는 보정하지 않는다.** 공격 애니메이션은 이미 클라 로컬 재생이고 판정만 + 서버다. 명중 판정에 공격자의 지연을 되돌리는 지연 보상은 안티사이트 면이 커서 + 별도 검토가 필요하다. + +### 검증 + +- 클라 테스트에 모의 RTT **20/150/250/400ms** 케이스를 추가했다. 250ms에서 보정된 + 클라와 측정값이 없는 클라가 같은 응답을 받는 순간, 전자는 서버의 현재 위치를, + 후자는 편도 지연만큼 과거 위치를 보여준다(3m/s에서 0.375m 차이). +- 회전·정지 중지를 걸친 입력 변경에서도 Approved 구간이 계속 재생되는지, + 정지·클릭·`finish`는 여전히 앵커를 버리는지 함께 검증한다. +- `LinkLatency` 단위 14개: 첫 응답 전에는 0, 2초 페이싱, 어긋난 재전송 seq 거부, + 재사용된 seq 거부, 스파이크 점진 수렴, reset 후 즉시 재탐색. +- **실행 중인 서버에 붙여 프로토콜을 확인했다.** 서버를 로컬에서 띄우고 배포된 + WASM 코덱으로 `ClientInfo` → `Ping` ×6을 보내 왕복을 확인했다. 양쪽 + `protocol_version` 104, 6/6 `Pong` 수신, 전부 seq와 echo `client_time_ms`가 일치했다. + 유닛 테스트가 못 잡는 것(실제 소켓, 실제 MessagePack 프레이밍, 실제 서버 바이너리가 + 새 변수를 같은 이름으로 agree하는지)을 이 검사가 덮는다. 루프백 RTT 평균 6.8ms. +- 위 실행 검사가 offset 신호를 무의미하다고 확인해 줬다(`server_time_ms`는 1.79e12 + epoch ms). 계측되지 않던 오프셋 코드가 그 measurement로 드러났고, 지워서 + 위 항목에 남겼다. +- **이 환경에는 C 링크러가 없어 못 돌렸던 검사들을 실제로 돌렸다:** + `cargo check --workspace --all-targets` 통과, `cargo test --workspace --locked` + **1,869개 통과·0 실패**, `cargo clippy --workspace --all-targets --locked -D warnings` + 통과, `cargo fmt --all --check` 통과, WASM 빌드 성공, 클라 **1,387개 통과** + (남은 3개 파일은 `tools/fetch-assets.sh`가 가져오는 3D 모델 부재). + ## 확인된 근거와 남은 검증 | 항목 | 현재 상태 | @@ -649,7 +723,10 @@ sequenceDiagram | 네트워크 왕복 | 현재 머신 WebSocket 평균 4.8ms, 외부 TCP 지점 서울 3–4ms·LA 131–184ms·뉴욕 201–209ms | | 예측 없이 출발 | 추가 서버 대기 20ms 가정 시 현재 머신 약 35ms·LA 161–215ms·뉴욕 231–240ms. 계산한 예상치 | | 예측 출발의 필요성 | 1단계의 실제 조작 반응을 보고 판단. 충분하면 2단계 생략 | +| 링크 지연 보정 | 프로토콜 104로 왕복 탐색 도입. 편도 지연만큼 앵커를 뒤로 dating하고, 입력 변경 시 승인 구간 유지 추가. 모의 RTT 20/150/250/400ms로 검증, A/B로 250ms에서 표시 지연 0.375m → 0.097m | +| 입력→첫 모션 지연 | **여전히 1 RTT.** 보정은 표시 시계만 건드리므로 250ms 회선에서 250ms 지연이 남는다. 2단계 예측 출발이 없이는 더 줄지 않는다 | | 예측 일치·교정 빈도 | 미측정. 선택적 2단계를 도입할 때만 경로와 진행 시각 차이를 검증 | +| 서버 빌드 검증 | sysroot를 직접 구성해 전부 실행. `cargo check`·`cargo test` 1,869/0 실패·`clippy -D warnings`·`cargo fmt --check` 통과. 실행 중인 서버에 WASM 코덱으로 Ping/Pong 6/6 왕복 확인 | | 동접 5,000명 전체 용량 | 미검증. 탐색 시험에는 게임 프로세스의 이동·전투·AI·잠금·AOI 전파가 포함되지 않음 | 아래는 최초 검증 계획이다. 구현 및 자동 검증 결과는 문서 상단에 기록한다. 실제 회선의 조작감으로 선택적 2단계를 판단한다. diff --git a/server/src/connection.rs b/server/src/connection.rs index 57c9f6c8..2cca6f05 100644 --- a/server/src/connection.rs +++ b/server/src/connection.rs @@ -1691,6 +1691,17 @@ async fn handle_client_message( state.last_heartbeat = std::time::Instant::now(); } + ClientMessage::Ping { + seq, + client_time_ms, + } => { + return Ok(vec![ServerMessage::Pong { + seq, + client_time_ms, + server_time_ms: GameState::now_ms(), + }]); + } + ClientMessage::EnvReport(r) => { if state.env_reported { return Ok(vec![]); diff --git a/shared/src/lib.rs b/shared/src/lib.rs index 83012e32..1198ff83 100644 --- a/shared/src/lib.rs +++ b/shared/src/lib.rs @@ -192,7 +192,9 @@ pub const NPC_TOKEN_FILENAME: &str = "npc_token"; /// v101: tip hats accept optional song requests delivered to the performer. /// v102: rain cells drift with the seasonal wind, change size and dry in the lee of ridges. /// v103: winter cells fall as snow, and WeatherSync can force snow. -pub const PROTOCOL_VERSION: u32 = 103; +/// v104: a Ping/Pong round-trip probe, so a client far from the region can +/// measure its own link instead of guessing at playback lag. +pub const PROTOCOL_VERSION: u32 = 104; /// Fingerprint of the dungeon layout generator this build compiled, stamped by /// `build.rs`. Layouts never travel the wire — both sides generate them from diff --git a/shared/src/messages.rs b/shared/src/messages.rs index 4002e9a7..10006835 100644 --- a/shared/src/messages.rs +++ b/shared/src/messages.rs @@ -906,6 +906,17 @@ pub enum ClientMessage { #[serde(default)] target_player_id: Option, }, + /// Round-trip probe, answered ahead of any queued work so it measures the + /// link rather than the server's load. `client_time_ms` rides back in + /// `ServerMessage::Pong` so the client needs neither clock to be trusted. + /// Subject to the pre-auth message budget like any other frame. + /// + /// Appended rather than placed beside `Heartbeat`: MessagePack encodes + /// these variants by index, so a new one belongs at the end. + Ping { + seq: u32, + client_time_ms: u64, + }, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -1884,6 +1895,16 @@ pub enum ServerMessage { monster_id: Option, remaining_ms: u64, }, + /// Answer to `ClientMessage::Ping`, sent straight down the same socket. + /// `server_time_ms` lets the client estimate its clock offset; that offset + /// only ever dates a probe, it never places the player. + /// + /// Appended for the same reason as `ClientMessage::Ping`. + Pong { + seq: u32, + client_time_ms: u64, + server_time_ms: u64, + }, } pub use crate::entity::PlayerId; @@ -2092,7 +2113,7 @@ impl ServerMessage { Self::GameTimeSync { .. } | Self::WeatherSync { .. } | Self::ServerNotice { .. } => { DeliveryClass::Global } - Self::WorldUpdate { .. } => DeliveryClass::Control, + Self::WorldUpdate { .. } | Self::Pong { .. } => DeliveryClass::Control, } } }