diff --git a/src/client/realm/service-connection-peer.ts b/src/client/realm/service-connection-peer.ts index 9b4f1e0..a8274af 100644 --- a/src/client/realm/service-connection-peer.ts +++ b/src/client/realm/service-connection-peer.ts @@ -157,7 +157,7 @@ export class RealmPeer extends SimplePeer { // explicit sync requested if (ping.dat.peerSync) { - const actions = await this.#sync.fetchSyncDelta(ping.dat.peerClocks) + const actions = await this.#sync.buildSyncDelta(ping.dat.peerClocks) if (actions.length) { this.sendJson(actions.map((a) => a.action)) } @@ -168,7 +168,7 @@ export class RealmPeer extends SimplePeer { // when a peer responds to a ping (lazy sync) #receivePong = async (pong: protocol.RealmRtcPongResponse) => { - const actions = await this.#sync.fetchSyncDelta(pong.dat.peerClocks) + const actions = await this.#sync.buildSyncDelta(pong.dat.peerClocks) if (actions.length) { this.sendJson(actions.map((a) => a.action)) } diff --git a/src/client/realm/service-connection-sync.ts b/src/client/realm/service-connection-sync.ts index 6e1a360..be80aa6 100644 --- a/src/client/realm/service-connection-sync.ts +++ b/src/client/realm/service-connection-sync.ts @@ -49,7 +49,7 @@ export class RealmSyncManager { return states } - async fetchSyncDelta(clocks: Record): Promise { + async buildSyncDelta(clocks: Record): Promise { const known = await this.#db.actions.orderBy('actor').uniqueKeys() const results = await Promise.all( known.map(async (actor) => { diff --git a/src/client/realm/service-connection.ts b/src/client/realm/service-connection.ts index 1bee48a..f5270e3 100644 --- a/src/client/realm/service-connection.ts +++ b/src/client/realm/service-connection.ts @@ -13,6 +13,7 @@ import {streamSocketJson, takeSocketJson} from '#common/socket' import {Database} from '#client/root/service-database.js' +import {sleep} from '#common/async/sleep.js' import {RealmPeer} from './service-connection-peer' import {RealmSyncManager} from './service-connection-sync' import {RealmIdentity} from './service-identity' @@ -35,6 +36,7 @@ export interface ConnectionOptions { export class RealmConnection extends EventTarget { #socket: WebSocket #connectopts: ConnectionOptions + #abort: AbortController #identities: Map #peers: Map @@ -43,6 +45,9 @@ export class RealmConnection extends EventTarget { #identity: RealmIdentity #sync: RealmSyncManager + #serverseq = 0 + #serversync = false + constructor( url: string, db: Database, @@ -51,6 +56,7 @@ export class RealmConnection extends EventTarget { options: ConnectionOptions, ) { super() + this.#abort = new AbortController() this.#identity = identity this.#connectopts = options @@ -110,9 +116,14 @@ export class RealmConnection extends EventTarget { console.debug('broadcasting:', self, data) const json = JSON.stringify(data) + this.#peers.forEach((peer, identid) => { if (self || identid !== this.#identity.identid) peer.send(json) }) + + if (this.#serversync) { + this.#socket.send(json) + } } destroy() { @@ -127,6 +138,9 @@ export class RealmConnection extends EventTarget { peer.destroy() } + // shutdown loops + this.#abort.abort() + this.#peers.clear() this.#nonces.clear() } @@ -221,11 +235,22 @@ export class RealmConnection extends EventTarget { // we initiate connections outbound when starting up for (const peerid of resp.dat.peers) { if (peerid === this.#identity.identid) continue + if (peerid === 'socket') { + this.#serversync = true + continue + } console.debug('connecting...:', peerid) this.#connectPeer(peerid, true) } + // initiate the socket loop if we want it + if (this.#serversync) { + this.#socketLoop().catch((err: unknown) => { + console.error('error in the socket sync loop!', err) + }) + } + // finally, we're connected this.#dispatchCustomEvent('wsauth', resp) } catch (exc) { @@ -246,10 +271,6 @@ export class RealmConnection extends EventTarget { } switch (parse.data.msg) { - case 'realm.rtc.pong': - console.debug('got a pong response', parse) - return - case 'realm.rtc.signal': { const jwt = await jwtPayload(protocol.realmRtcSignalPayloadSchema).parseAsync( parse.data.dat.signed, @@ -300,6 +321,17 @@ export class RealmConnection extends EventTarget { this.#disconnectPeer(parse.data.dat.identid) this.#dispatchCustomEvent('peerleft', {identid: parse.data.dat.identid}) return + + case 'realm.rtc.pong': { + console.debug('got a pong response from the server', parse) + + const actions = await this.#sync.buildSyncDelta(parse.data.dat.peerClocks) + if (actions.length) { + this.#socket.send(JSON.stringify(actions.map((a) => a.action))) + } + + return + } } } @@ -313,6 +345,26 @@ export class RealmConnection extends EventTarget { this.destroy() } + #socketLoop = async () => { + this.#abort.signal.throwIfAborted() + + await this.#pingSocket() + while (!this.#abort.signal.aborted) { + await sleep(30_000) + await this.#pingSocket() + } + } + + async #pingSocket() { + const peerClocks = await this.#sync.buildSyncState() + this.#socketSend({ + typ: 'req', + msg: 'realm.rtc.ping', + seq: this.#serverseq++, + dat: {peerClocks, peerSync: true}, + }) + } + // peers #connectPeer(remoteid: IdentID, initiator: boolean): RealmPeer { diff --git a/src/client/skypod/context.tsx b/src/client/skypod/context.tsx index 3c91f5d..736c5ec 100644 --- a/src/client/skypod/context.tsx +++ b/src/client/skypod/context.tsx @@ -190,8 +190,10 @@ export const SkypodProvider: preact.FunctionComponent<{children: preact.Componen } connection.addEventListener('peerdata', handler as EventListener) + connection.addEventListener('wsdata', handler as EventListener) return () => { connection.removeEventListener('peerdata', handler as EventListener) + connection.removeEventListener('wsdata', handler as EventListener) } }, [context, clock, realm.value]) diff --git a/src/common/protocol/messages.ts b/src/common/protocol/messages.ts index 2b1bfb7..91b5253 100644 --- a/src/common/protocol/messages.ts +++ b/src/common/protocol/messages.ts @@ -10,6 +10,9 @@ import { makeResponseSchema, } from './schema' +export const socketPeerIdSchema = z.literal('socket') +export type SocketPeerId = z.infer + /// preauth export const preauthRegisterReqSchema = makeRequestSchema( @@ -30,7 +33,7 @@ export const preauthExchangeInviteReqSchema = makeRequestSchema( export const preauthRespSchema = makeResponseSchema( 'preauth.authn', z.object({ - peers: z.array(IdentBrand.schema), + peers: z.array(z.union([IdentBrand.schema, socketPeerIdSchema])), identities: z.record(z.string(), jwkSchema), }), ) diff --git a/src/common/strict-map.ts b/src/common/strict-map.ts index 9ef1985..28f01ce 100644 --- a/src/common/strict-map.ts +++ b/src/common/strict-map.ts @@ -22,6 +22,21 @@ export class StrictMap extends Map { return this.get(key)! } + /** get a value from the map, creating it (async) if not present */ + async ensureAsync(key: K, maker: () => Promise): Promise { + if (!this.has(key)) { + const value = await maker() + if (!this.has(key)) { + this.set(key, value) + } else { + console.warn('ensureasync was out-raced!', this, key) + } + } + + // eslint-disable-next-line @typescript-eslint/no-non-null-assertion + return this.get(key)! + } + /** update a value in the map, removing if undefined is returned */ update(key: K, update: (prev?: V) => V | undefined): void { const prev = this.get(key) diff --git a/src/server/realm-storage.ts b/src/server/realm-storage.ts new file mode 100644 index 0000000..b193fb2 --- /dev/null +++ b/src/server/realm-storage.ts @@ -0,0 +1,62 @@ +import {z} from 'zod/v4' + +import {IdentBrand, IdentID, RealmID} from '#common/protocol' +import {LCTimestamp, LogicalClock} from '#common/protocol/logical-clock' +import {actionMessageSchema} from '#common/protocol/schema' +import {StrictMap} from '#common/strict-map' + +export type IncomingAction = z.infer +export type StoredAction = z.infer +export const storedActionSchema = z.object({ + actor: IdentBrand.schema, + clock: LogicalClock.schema, + action: z.unknown(), +}) + +export class RealmStorage { + static async list(): Promise { + await Promise.resolve() + return [] + } + + static async ensure( + realmid: RealmID, + registrantid: IdentID, + registrantkey: CryptoKey, + ): Promise { + const storage = new RealmStorage(realmid) + await storage.admitIdentity(registrantid, registrantkey) + + // todo, look it up in a directory + + return storage + } + + realmid: RealmID + identities: StrictMap + + private constructor(realmid: RealmID) { + this.realmid = realmid + this.identities = new StrictMap() + } + + async admitIdentity(inviteeid: IdentID, inviteekey: CryptoKey) { + await Promise.resolve() + this.identities.set(inviteeid, inviteekey) + } + + async handleActions(actions: IncomingAction[]) { + await Promise.resolve() + console.log('incoming actions', actions) + } + + async buildSyncState(): Promise> { + await Promise.resolve() + return {} + } + + async buildSyncDelta(state: Record): Promise { + await Promise.resolve() + return [] + } +} diff --git a/src/server/routes-socket/handler-preauth.ts b/src/server/routes-socket/handler-preauth.ts index aac26a3..e959397 100644 --- a/src/server/routes-socket/handler-preauth.ts +++ b/src/server/routes-socket/handler-preauth.ts @@ -79,7 +79,7 @@ async function preauthValidateInvitation(realmid: RealmID, invitation: JWTToken) throw new Error('invitation already used!') const inviterid = IdentBrand.parse(invitation.claims.iss) - const inviterkey = realm.identities.require(inviterid) + const inviterkey = realm.storage.identities.require(inviterid) await verifyJwtToken(invitation.token, inviterkey, {subject: 'invitation'}) } catch (exc) { const err = normalizeError(exc) @@ -95,7 +95,7 @@ async function authenticatePreauth( try { const realmMap = await realms.ensureRealmMap() const realm = realmMap.require(realmid) - const pubkey = realm.identities.require(identid) + const pubkey = realm.storage.identities.require(identid) // at this point we no langer care about the payload // but this throws as a side-effect if the token is invalid @@ -111,11 +111,12 @@ async function preauthResponse( auth: realms.AuthenticatedIdentity, seq?: number, ): Promise { - const peers = Array.from(auth.realm.sockets.keys()) - const identities: Record = {} + const peers: Array = Array.from(auth.realm.sockets.keys()) + peers.unshift('socket') - for (const identid of auth.realm.identities.keys()) { - const pubkey = auth.realm.identities.require(identid) + const identities: Record = {} + for (const identid of auth.realm.storage.identities.keys()) { + const pubkey = auth.realm.storage.identities.require(identid) identities[identid] = await jwkExport.parseAsync(pubkey) } diff --git a/src/server/routes-socket/handler-realm.ts b/src/server/routes-socket/handler-realm.ts index 42f95fa..a55d88c 100644 --- a/src/server/routes-socket/handler-realm.ts +++ b/src/server/routes-socket/handler-realm.ts @@ -1,16 +1,20 @@ import {WebSocket} from 'isomorphic-ws' import {z} from 'zod/v4' -import {jwkExport} from '#common/crypto/jwks.js' +import {jwkExport} from '#common/crypto/jwks' import {normalizeProtocolError, ProtocolError} from '#common/errors' import * as protocol from '#common/protocol' +import {actionMessageSchema} from '#common/protocol/schema' import {streamSocket} from '#common/socket' + import * as realm from '#server/routes-socket/state' // what can the server handle? const incomingMessageSchema = z.union([ protocol.realmBroadcastEventSchema, protocol.realmRtcSignalEventSchema, + protocol.realmRtcPingPongMessageSchema, + z.array(actionMessageSchema), ]) /** @@ -29,6 +33,11 @@ export async function realmHandler( for await (const msg of streamSocket(ws, {signal})) { try { const data = await incomingParser.parseAsync(msg) + if (Array.isArray(data)) { + await auth.realm.storage.handleActions(data) + continue + } + switch (data.msg) { case 'realm.broadcast': realmBroadcast(auth, data.dat, data.dat.recipients) @@ -38,6 +47,14 @@ export async function realmHandler( realmBroadcast(auth, data, [data.dat.remoteid]) continue + case 'realm.rtc.ping': + await socketPeerPing(ws, auth, data) + continue + + case 'realm.rtc.pong': + await socketPeerPong(ws, auth, data) + continue + default: console.error('unknown message!', msg) throw new ProtocolError(`unknown message type!`, 400) @@ -92,7 +109,9 @@ function realmBroadcast( recipients: protocol.IdentID[] | boolean = false, ) { const echo = recipients === true || Array.isArray(recipients) - const recips = Array.isArray(recipients) ? recipients : Array.from(auth.realm.identities.keys()) + const recips = Array.isArray(recipients) + ? recipients + : Array.from(auth.realm.storage.identities.keys()) const json = JSON.stringify(payload) for (const recip of recips) { @@ -106,3 +125,38 @@ function realmBroadcast( } } } + +async function socketPeerPing( + ws: WebSocket, + auth: realm.AuthenticatedIdentity, + ping: protocol.RealmRtcPingRequest, +) { + const peerClocks = await auth.realm.storage.buildSyncState() + const response: protocol.RealmRtcPongResponse = { + typ: 'res', + msg: 'realm.rtc.pong', + seq: ping.seq, + dat: {peerClocks}, + } + ws.send(JSON.stringify(response)) + + if (ping.dat.peerSync) { + const actions = await auth.realm.storage.buildSyncDelta(ping.dat.peerClocks) + if (actions.length) { + const actionsJson = actions.map((a) => a.action) + ws.send(JSON.stringify(actionsJson)) + } + } +} + +async function socketPeerPong( + ws: WebSocket, + auth: realm.AuthenticatedIdentity, + pong: protocol.RealmRtcPongResponse, +) { + const actions = await auth.realm.storage.buildSyncDelta(pong.dat.peerClocks) + if (actions.length) { + const actionsJson = actions.map((a) => a.action) + ws.send(JSON.stringify(actionsJson)) + } +} diff --git a/src/server/routes-socket/state.ts b/src/server/routes-socket/state.ts index ec61bdb..6846cda 100644 --- a/src/server/routes-socket/state.ts +++ b/src/server/routes-socket/state.ts @@ -3,6 +3,8 @@ import WebSocket from 'isomorphic-ws' import {IdentID, RealmID} from '#common/protocol' import {StrictMap} from '#common/strict-map' +import {RealmStorage} from '#server/realm-storage' + /** An authenticated identity; only handed out in response to successful authentication. */ export interface AuthenticatedIdentity { realm: Realm @@ -13,9 +15,9 @@ export interface AuthenticatedIdentity { export interface Realm { realmid: RealmID + storage: RealmStorage nonces: Set sockets: StrictMap - identities: StrictMap } // async interface, because we're going to put this in storage at some point @@ -47,16 +49,15 @@ export function ensureRegisteredRealm( registrantid: IdentID, registrantkey: CryptoKey, ): Promise { - const realm = realmMap.ensure(realmid, () => ({ - realmid, - nonces: new Set(), - sockets: new StrictMap(), - identities: new StrictMap([[registrantid, registrantkey]]), - })) - - // hack for now, allow any registration to work - realm.identities.ensure(registrantid, () => registrantkey) - return Promise.resolve(realm) + return realmMap.ensureAsync(realmid, async () => { + const storage = await RealmStorage.ensure(realmid, registrantid, registrantkey) + return { + realmid, + storage, + nonces: new Set(), + sockets: new StrictMap(), + } + }) } export function validateNonce(realmid: RealmID, nonce: string): Promise { @@ -71,9 +72,7 @@ export function validateNonce(realmid: RealmID, nonce: string): Promise export function admitToRealm(realmid: RealmID, inviteeid: IdentID, inviteekey: CryptoKey) { const realm = realmMap.require(realmid) - realm.identities.set(inviteeid, inviteekey) - - return Promise.resolve() + return realm.storage.admitIdentity(inviteeid, inviteekey) } export function attachSocket(realm: Realm, ident: IdentID, socket: WebSocket) {