diff --git a/examples/nuxtjs-live-avatars/package-lock.json b/examples/nuxtjs-live-avatars/package-lock.json index f062bf6830b..334fb9e0957 100644 --- a/examples/nuxtjs-live-avatars/package-lock.json +++ b/examples/nuxtjs-live-avatars/package-lock.json @@ -7188,7 +7188,9 @@ "license": "MIT" }, "node_modules/node-forge": { - "version": "1.3.3", + "version": "1.4.0", + "resolved": "https://registry.npmjs.org/node-forge/-/node-forge-1.4.0.tgz", + "integrity": "sha512-LarFH0+6VfriEhqMMcLX2F7SwSXeWwnEAJEsYm5QKWchiVYVvJyV9v7UDvUv+w5HO23ZpQTXDv/GxdDdMyOuoQ==", "dev": true, "license": "(BSD-3-Clause OR GPL-2.0)", "engines": { diff --git a/packages/liveblocks-server/src/Room.ts b/packages/liveblocks-server/src/Room.ts index 594bedf3de4..72c0a94922a 100644 --- a/packages/liveblocks-server/src/Room.ts +++ b/packages/liveblocks-server/src/Room.ts @@ -28,7 +28,7 @@ import { nodeStreamToCompactNodes, OpCode, raise, - ServerMsgCode, + ServerMsgCode as CoreServerMsgCode, tryParseJson, WebsocketCloseCodes as CloseCode, } from "@liveblocks/core"; @@ -39,27 +39,48 @@ import { nanoid } from "nanoid"; import type { Guid } from "~/decoders"; import { clientMsgDecoder } from "~/decoders"; -import type { IServerWebSocket, IStorageDriver } from "~/interfaces"; +import type { + IServerWebSocket, + IStorageDriver, + ListFeedMessagesOptions, + ListFeedMessagesResult, + ListFeedsOptions, + ListFeedsResult, +} from "~/interfaces"; import { Logger } from "~/lib/Logger"; import { makeNewInMemoryDriver } from "~/plugins/InMemoryDriver"; import type { + AddFeedClientMsg, + AddFeedMessageClientMsg, ClientMsg as GenericClientMsg, + DeleteFeedClientMsg, + DeleteFeedMessageClientMsg, + FeedMessagesServerMsg, + FeedsServerMsg, + FetchFeedMessagesClientMsg, + FetchFeedsClientMsg, IgnoredOp, Op, ServerMsg as GenericServerMsg, ServerWireOp, + UpdateFeedClientMsg, + UpdateFeedMessageClientMsg, } from "~/protocol"; -import { ProtocolVersion } from "~/protocol"; import { Storage } from "~/Storage"; import { YjsStorage } from "~/YjsStorage"; import { tryCatch } from "./lib/tryCatch"; import { UniqueMap } from "./lib/UniqueMap"; -import type { LeasedSession } from "./types"; +import { feedFailureServerMsg, feedRequestFailed } from "./protocol/feedErrors"; +import { FeedMsgCode, FeedRequestErrorCode } from "./protocol/feedMessages"; +import { ProtocolVersion } from "./protocol/ProtocolVersion"; +import type { Feed, FeedMessage, LeasedSession } from "./types"; import { isLeasedSessionExpired, makeRoomStateMsg } from "./utils"; const messagesDecoder = array(clientMsgDecoder); +// Temporary patch +const ServerMsgCode = { ...CoreServerMsgCode, ...FeedMsgCode }; const HIGHEST_PROTOCOL_VERSION = Math.max( ...Object.values(ProtocolVersion).filter( (v): v is number => typeof v === "number" @@ -89,8 +110,20 @@ export type Millis = Brand; export type SessionKey = Brand; export type PreSerializedServerMsg = Brand; -type ClientMsg = GenericClientMsg; -type ServerMsg = GenericServerMsg; +type ClientMsg = + | GenericClientMsg + | FetchFeedsClientMsg + | FetchFeedMessagesClientMsg + | AddFeedClientMsg + | UpdateFeedClientMsg + | DeleteFeedClientMsg + | AddFeedMessageClientMsg + | UpdateFeedMessageClientMsg + | DeleteFeedMessageClientMsg; +type ServerMsg = + | GenericServerMsg + | FeedMessagesServerMsg + | FeedsServerMsg; // Temp patched server msgs /** * Creates a collector for deferred promises (side effects that should run @@ -1189,6 +1222,117 @@ export class Room { } } + // --------------------------------------------------------------------------- + // Feed APIs + // --------------------------------------------------------------------------- + + /** + * List feeds with pagination and filtering. + */ + public async listFeeds(options?: ListFeedsOptions): Promise { + return await this.driver.list_feeds(options); + } + + /** + * Get a specific feed by feed ID. + */ + public async getFeed(feedId: string): Promise { + return await this.driver.get_feed(feedId); + } + + /** + * Create a new feed. + * If timestamp is not provided, current server time is used. + */ + public async createFeed( + feed: Omit & { timestamp?: number } + ): Promise { + const now = feed.timestamp ?? Date.now(); + const fullFeed: Feed = { + ...feed, + createdAt: now, + updatedAt: now, + }; + await this.driver.create_feed(fullFeed); + return fullFeed; + } + + /** + * Update a feed's metadata. + */ + public async updateFeedMetadata( + feedId: string, + metadata: Json + ): Promise { + await this.driver.update_feed_metadata(feedId, metadata); + } + + /** + * Delete a feed. + */ + public async deleteFeed(feedId: string): Promise { + await this.driver.delete_feed(feedId); + } + + /** + * List feed messages for a feed with pagination. + */ + public async listFeedMessages( + feedId: string, + options?: ListFeedMessagesOptions + ): Promise { + return await this.driver.list_feed_messages(feedId, options); + } + + /** + * Add a message to a feed. + * If message id is not provided, a unique ID is automatically generated. + * If timestamp is not provided, current server time is used. + */ + public async addFeedMessage( + feedId: string, + message: Omit & + Partial> & { timestamp?: number } + ): Promise { + const now = message.timestamp ?? Date.now(); + const fullMessage: FeedMessage = { + id: message.id ?? nanoid(), + createdAt: now, + updatedAt: now, + data: message.data, + }; + await this.driver.add_feed_message(feedId, fullMessage); + return fullMessage; + } + + /** + * Update a feed message's data. + * Returns the updated message. + */ + public async updateFeedMessage( + feedId: string, + messageId: string, + data: Json, + timestamp?: number + ): Promise { + return await this.driver.update_feed_message( + feedId, + messageId, + data, + timestamp ?? Date.now() + ); + } + + /** + * Delete a feed message. + */ + public async deleteFeedMessage( + feedId: string, + messageId: string + ): Promise { + await this.driver.delete_feed_message(feedId, messageId); + } + /** * Will send the given ServerMsg through all Session, except the Session * where the message originates from. @@ -1647,6 +1791,185 @@ export class Room { break; } + // Feed messages + case FeedMsgCode.FETCH_FEEDS: { + const fetchMsg = msg; + const [result, err] = await tryCatch( + this.listFeeds({ + cursor: fetchMsg.cursor, + since: fetchMsg.since, + limit: fetchMsg.limit, + metadata: fetchMsg.metadata, + }) + ); + if (err) { + replyImmediately(feedFailureServerMsg(fetchMsg.requestId, err)); + break; + } + replyImmediately({ + type: FeedMsgCode.FEEDS_LIST, + requestId: fetchMsg.requestId, + feeds: result.feeds, + nextCursor: result.nextCursor, + }); + break; + } + + case FeedMsgCode.FETCH_FEED_MESSAGES: { + const fetchMsg = msg; + const [result, err] = await tryCatch( + this.listFeedMessages(fetchMsg.feedId, { + cursor: fetchMsg.cursor, + since: fetchMsg.since, + limit: fetchMsg.limit, + }) + ); + if (err) { + replyImmediately(feedFailureServerMsg(fetchMsg.requestId, err)); + break; + } + replyImmediately({ + type: FeedMsgCode.FEED_MESSAGES_LIST, + requestId: fetchMsg.requestId, + feedId: fetchMsg.feedId, + messages: result.messages, + nextCursor: result.nextCursor, + }); + break; + } + + case FeedMsgCode.ADD_FEED: { + const addMsg = msg; + const [feed, err] = await tryCatch( + this.createFeed({ + feedId: addMsg.feedId, + metadata: (addMsg.metadata as Json) ?? {}, + timestamp: addMsg.timestamp, + }) + ); + if (err) { + replyImmediately(feedFailureServerMsg(addMsg.requestId, err)); + break; + } + const feedsMsg = { + type: FeedMsgCode.FEEDS_ADDED, + feeds: [feed], + }; + replyImmediately(feedsMsg); + scheduleFanOut(feedsMsg); + break; + } + + case FeedMsgCode.UPDATE_FEED: { + const updateMsg = msg; + const [, metaErr] = await tryCatch( + this.updateFeedMetadata(updateMsg.feedId, updateMsg.metadata as Json) + ); + if (metaErr) { + replyImmediately(feedFailureServerMsg(updateMsg.requestId, metaErr)); + break; + } + const feed = await this.getFeed(updateMsg.feedId); + if (!feed) { + replyImmediately( + feedRequestFailed( + updateMsg.requestId, + FeedRequestErrorCode.FEED_NOT_FOUND + ) + ); + break; + } + const feedsMsg = { + type: FeedMsgCode.FEEDS_UPDATED, + feeds: [feed], + }; + replyImmediately(feedsMsg); + scheduleFanOut(feedsMsg); + break; + } + + case FeedMsgCode.DELETE_FEED: { + const deleteMsg = msg; + const [, err] = await tryCatch(this.deleteFeed(deleteMsg.feedId)); + if (err) { + replyImmediately(feedFailureServerMsg(deleteMsg.requestId, err)); + break; + } + const feedDeletedMsg = { + type: FeedMsgCode.FEED_DELETED, + feedId: deleteMsg.feedId, + }; + replyImmediately(feedDeletedMsg); + scheduleFanOut(feedDeletedMsg); + break; + } + + case FeedMsgCode.ADD_FEED_MESSAGE: { + const addMsg = msg; + const [message, err] = await tryCatch( + this.addFeedMessage(addMsg.feedId, { + data: addMsg.data as Json, + id: addMsg.id, + timestamp: addMsg.timestamp, + }) + ); + if (err) { + replyImmediately(feedFailureServerMsg(addMsg.requestId, err)); + break; + } + const feedMessagesMsg = { + type: FeedMsgCode.FEED_MESSAGES_ADDED, + feedId: addMsg.feedId, + messages: [message], + }; + replyImmediately(feedMessagesMsg); + scheduleFanOut(feedMessagesMsg); + break; + } + + case FeedMsgCode.UPDATE_FEED_MESSAGE: { + const updateMsg = msg; + const [message, err] = await tryCatch( + this.updateFeedMessage( + updateMsg.feedId, + updateMsg.messageId, + updateMsg.data as Json, + updateMsg.timestamp + ) + ); + if (err) { + replyImmediately(feedFailureServerMsg(updateMsg.requestId, err)); + break; + } + const feedMessagesMsg = { + type: FeedMsgCode.FEED_MESSAGES_UPDATED, + feedId: updateMsg.feedId, + messages: [message], + }; + replyImmediately(feedMessagesMsg); + scheduleFanOut(feedMessagesMsg); + break; + } + + case FeedMsgCode.DELETE_FEED_MESSAGE: { + const deleteMsg = msg; + const [, err] = await tryCatch( + this.deleteFeedMessage(deleteMsg.feedId, deleteMsg.messageId) + ); + if (err) { + replyImmediately(feedFailureServerMsg(deleteMsg.requestId, err)); + break; + } + const feedMessagesMsg = { + type: FeedMsgCode.FEED_MESSAGES_DELETED, + feedId: deleteMsg.feedId, + messageIds: [deleteMsg.messageId], + }; + replyImmediately(feedMessagesMsg); + scheduleFanOut(feedMessagesMsg); + break; + } + default: { try { return assertNever(msg, "Unrecognized client msg"); diff --git a/packages/liveblocks-server/src/decoders/ClientMsg.ts b/packages/liveblocks-server/src/decoders/ClientMsg.ts index 8f83d56b94d..ebef32a2ca2 100644 --- a/packages/liveblocks-server/src/decoders/ClientMsg.ts +++ b/packages/liveblocks-server/src/decoders/ClientMsg.ts @@ -22,6 +22,8 @@ import { array, boolean, constant, + jsonObject, + nonEmptyString, number, object, optional, @@ -30,15 +32,29 @@ import { } from "decoders"; import type { + AddFeedClientMsg, + AddFeedMessageClientMsg, BroadcastEventClientMsg, ClientMsg, + DeleteFeedClientMsg, + DeleteFeedMessageClientMsg, + FetchFeedMessagesClientMsg, + FetchFeedsClientMsg, FetchStorageClientMsg, FetchYDocClientMsg, + UpdateFeedClientMsg, + UpdateFeedMessageClientMsg, UpdatePresenceClientMsg, UpdateStorageClientMsg, UpdateYDocClientMsg, } from "~/protocol"; +import { FeedMsgCode } from "~/protocol"; +import { + feedMetadataUpdateDecoder, + fetchFeedsMetadataFilterDecoder, + optionalFeedMetadataDecoder, +} from "./feedMetadata"; import { jsonObjectYolo, jsonYolo } from "./jsonYolo"; import { op } from "./Op"; import type { YUpdate, YVector } from "./y-types"; @@ -79,6 +95,71 @@ const updateYDocClientMsg: Decoder = object({ v2: optional(boolean), }); +// Feed message decoders +const fetchFeedsClientMsg: Decoder = object({ + type: constant(FeedMsgCode.FETCH_FEEDS), + requestId: string, + cursor: optional(string), + since: optional(number), + limit: optional(number), + metadata: fetchFeedsMetadataFilterDecoder, +}); + +const fetchFeedMessagesClientMsg: Decoder = object({ + type: constant(FeedMsgCode.FETCH_FEED_MESSAGES), + requestId: string, + feedId: nonEmptyString, + cursor: optional(string), + since: optional(number), + limit: optional(number), +}); + +const addFeedClientMsg: Decoder = object({ + type: constant(FeedMsgCode.ADD_FEED), + feedId: string, + metadata: optionalFeedMetadataDecoder, + timestamp: optional(number), + requestId: optional(string), +}); + +const updateFeedClientMsg: Decoder = object({ + type: constant(FeedMsgCode.UPDATE_FEED), + feedId: string, + metadata: feedMetadataUpdateDecoder, + requestId: optional(string), +}); + +const deleteFeedClientMsg: Decoder = object({ + type: constant(FeedMsgCode.DELETE_FEED), + feedId: string, + requestId: optional(string), +}); + +const addFeedMessageClientMsg: Decoder = object({ + type: constant(FeedMsgCode.ADD_FEED_MESSAGE), + feedId: string, + data: jsonObject, + id: optional(string), + timestamp: optional(number), + requestId: optional(string), +}); + +const updateFeedMessageClientMsg: Decoder = object({ + type: constant(FeedMsgCode.UPDATE_FEED_MESSAGE), + feedId: string, + messageId: string, + data: jsonObject, + timestamp: optional(number), + requestId: optional(string), +}); + +const deleteFeedMessageClientMsg: Decoder = object({ + type: constant(FeedMsgCode.DELETE_FEED_MESSAGE), + feedId: string, + messageId: string, + requestId: optional(string), +}); + export const clientMsgDecoder: Decoder> = taggedUnion("type", { [ClientMsgCode.UPDATE_PRESENCE]: updatePresenceClientMsg, @@ -87,7 +168,18 @@ export const clientMsgDecoder: Decoder> = [ClientMsgCode.UPDATE_STORAGE]: updateStorageClientMsg, [ClientMsgCode.FETCH_YDOC]: fetchYDocClientMsg, [ClientMsgCode.UPDATE_YDOC]: updateYDocClientMsg, - }).describe("Must be a valid client message"); + [FeedMsgCode.FETCH_FEEDS]: fetchFeedsClientMsg, + [FeedMsgCode.FETCH_FEED_MESSAGES]: fetchFeedMessagesClientMsg, + [FeedMsgCode.ADD_FEED]: addFeedClientMsg, + [FeedMsgCode.UPDATE_FEED]: updateFeedClientMsg, + [FeedMsgCode.DELETE_FEED]: deleteFeedClientMsg, + [FeedMsgCode.ADD_FEED_MESSAGE]: addFeedMessageClientMsg, + [FeedMsgCode.UPDATE_FEED_MESSAGE]: updateFeedMessageClientMsg, + [FeedMsgCode.DELETE_FEED_MESSAGE]: deleteFeedMessageClientMsg, + } as unknown as Record< + string | number, + Decoder> + >).describe("Must be a valid client message"); export const transientClientMsgDecoder: Decoder> = taggedUnion("type", { diff --git a/packages/liveblocks-server/src/decoders/feedMetadata.ts b/packages/liveblocks-server/src/decoders/feedMetadata.ts new file mode 100644 index 00000000000..5067813b9bb --- /dev/null +++ b/packages/liveblocks-server/src/decoders/feedMetadata.ts @@ -0,0 +1,102 @@ +/** + * Copyright (c) Liveblocks Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +/** + * Feed metadata decoders — same rules as room/thread metadata (see @shared/common + * createMetadataDecoder / updateRoomMetadataDecoder). + */ + +import type { Decoder } from "decoders"; +import { + array, + either, + identifier, + nullable, + optional, + poja, + record, + select, + sized, + string, +} from "decoders"; + +/** Same cap as createMetadataDecoder in @shared/common */ +const MAX_METADATA_COUNT = 50; +const MAX_METADATA_VALUE_LIST_LENGTH = 50; + +export const feedMetadataIdDecoder = sized(identifier, { min: 1, max: 40 }); + +const metadataStringValue = sized(string, { max: 256 }); + +function createRoomMetadataValueDecoder(): Decoder { + return select( + either(string, poja).describe("Must be string or string[]"), + (x) => + typeof x === "string" + ? metadataStringValue + : array(metadataStringValue).refine( + (value) => value.length <= MAX_METADATA_VALUE_LIST_LENGTH, + `Must be at most ${MAX_METADATA_VALUE_LIST_LENGTH} items` + ) + ); +} + +const roomMetadataValueDecoder = createRoomMetadataValueDecoder(); + +/** + * A single metadata value: string, string[], or null (null clears a key on update). + */ +const feedMetadataNullableValueDecoder = nullable(roomMetadataValueDecoder); + +const feedMetadataRecordForCreate = record( + feedMetadataIdDecoder, + roomMetadataValueDecoder +).refine( + (value) => Object.keys(value).length <= MAX_METADATA_COUNT, + `Must be at most ${MAX_METADATA_COUNT} items` +); + +/** + * Optional feed metadata on create (WebSocket ADD_FEED, HTTP POST …/feed). + * Same rules as createMetadataDecoder / room metadata on create (no null values). + */ +export const optionalFeedMetadataDecoder = optional(feedMetadataRecordForCreate); + +/** + * Full metadata object for update (WebSocket UPDATE_FEED, HTTP PATCH …/feeds/:id). + * Same shape as updateRoomMetadataDecoder. + */ +export const feedMetadataUpdateDecoder = record( + feedMetadataIdDecoder, + feedMetadataNullableValueDecoder +); + +const feedMetadataRecordForFilter = record( + feedMetadataIdDecoder, + feedMetadataNullableValueDecoder +).refine( + (value) => Object.keys(value).length <= MAX_METADATA_COUNT, + `Must be at most ${MAX_METADATA_COUNT} items` +); + +/** + * Optional metadata filter for FETCH_FEEDS — same key/value rules as + * {@link feedMetadataUpdateDecoder} / updateRoomMetadataDecoder in @shared/common. + */ +export const fetchFeedsMetadataFilterDecoder = optional( + feedMetadataRecordForFilter +); diff --git a/packages/liveblocks-server/src/decoders/index.ts b/packages/liveblocks-server/src/decoders/index.ts index 0f1e9729613..2f421958ba2 100644 --- a/packages/liveblocks-server/src/decoders/index.ts +++ b/packages/liveblocks-server/src/decoders/index.ts @@ -16,5 +16,6 @@ */ export * from "./ClientMsg"; +export * from "./feedMetadata"; export * from "./jsonYolo"; export * from "./y-types"; diff --git a/packages/liveblocks-server/src/index.ts b/packages/liveblocks-server/src/index.ts index f06af8a9175..dc7be8d4a17 100644 --- a/packages/liveblocks-server/src/index.ts +++ b/packages/liveblocks-server/src/index.ts @@ -35,6 +35,8 @@ export { makeMetadataDB } from "~/MetadataDB"; export * from "~/protocol"; export * from "~/Room"; export type { + Feed, + FeedMessage, LeasedSession, NodeMap, NodeStream, diff --git a/packages/liveblocks-server/src/interfaces/IStorageDriver.ts b/packages/liveblocks-server/src/interfaces/IStorageDriver.ts index ab8792a892d..8c3358228be 100644 --- a/packages/liveblocks-server/src/interfaces/IStorageDriver.ts +++ b/packages/liveblocks-server/src/interfaces/IStorageDriver.ts @@ -28,7 +28,42 @@ import type { import type { YDocId } from "~/decoders/y-types"; import type { Logger } from "~/lib/Logger"; -import type { LeasedSession, Pos } from "~/types"; +import type { Feed, FeedMessage, LeasedSession, Pos } from "~/types"; + +/** + * Options for listing feeds with pagination and filtering. + */ +export type ListFeedsOptions = { + cursor?: string; // Pagination cursor + since?: number; // Unix timestamp in milliseconds - return feeds created after this time + limit?: number; // Max items to return (1-100) + metadata?: Record; // Metadata filters (key-value pairs, supports string/number/boolean/null) +}; + +/** + * Options for listing feed messages with pagination. + */ +export type ListFeedMessagesOptions = { + cursor?: string; // Pagination cursor + since?: number; // Unix timestamp in milliseconds - return messages created after this time + limit?: number; // Max items to return (1-100) +}; + +/** + * Result of listing feeds with pagination info. + */ +export type ListFeedsResult = { + feeds: Feed[]; + nextCursor?: string; // Cursor for next page, undefined if no more pages +}; + +/** + * Result of listing feed messages with pagination info. + */ +export type ListFeedMessagesResult = { + messages: FeedMessage[]; + nextCursor?: string; // Cursor for next page, undefined if no more pages +}; /** * An isolated, read-only copy of the storage document at a point in time. @@ -301,4 +336,77 @@ export interface IStorageDriver { * and reset the counter. */ takeRowsWritten?(): number; + + // --------------------------------------------------------------------------- + // Feed APIs + // --------------------------------------------------------------------------- + + /** + * List feeds with pagination, filtering, and metadata querying. + * Feeds are sorted by createdAt descending (newest first). + */ + list_feeds(options?: ListFeedsOptions): Awaitable; + + /** + * Get a specific feed by feed ID. + * Returns feed metadata only (without messages). + * Use list_feed_messages to retrieve messages for this feed. + * Returns undefined if the feed doesn't exist. + */ + get_feed(feedId: string): Awaitable; + + /** + * Create a new feed. + * If feedId already exists, throws an error. + */ + create_feed(feed: Feed): Awaitable; + + /** + * Update a feed's metadata. + * The feed must exist, otherwise throws an error. + */ + update_feed_metadata(feedId: string, metadata: Json): Awaitable; + + /** + * Delete a feed by feed ID. + * Also deletes all messages associated with the feed (via CASCADE). + * No-op if feed doesn't exist. + */ + delete_feed(feedId: string): Awaitable; + + /** + * List feed messages for a feed with pagination. + * Messages are sorted by createdAt descending (newest first). + */ + list_feed_messages( + feedId: string, + options?: ListFeedMessagesOptions + ): Awaitable; + + /** + * Add a message to a feed. + * The message must have id, createdAt, and updatedAt already set (handled by Room layer). + * The feed must exist, otherwise throws an error. + */ + add_feed_message(feedId: string, message: FeedMessage): Awaitable; + + /** + * Update a feed message's data. + * Returns the updated message. + * The feed and message must exist, otherwise throws an error. + * If timestamp is not provided, current server time is used. + * Messages are only updated if the provided timestamp is greater than or equal to the stored updatedAt. + */ + update_feed_message( + feedId: string, + messageId: string, + data: Json, + timestamp?: number + ): Awaitable; + + /** + * Delete a feed message. + * The feed and message must exist, otherwise throws an error. + */ + delete_feed_message(feedId: string, messageId: string): Awaitable; } diff --git a/packages/liveblocks-server/src/interfaces/index.ts b/packages/liveblocks-server/src/interfaces/index.ts index 1ad0ed11312..a61985b9d81 100644 --- a/packages/liveblocks-server/src/interfaces/index.ts +++ b/packages/liveblocks-server/src/interfaces/index.ts @@ -20,5 +20,9 @@ export type { IReadableSnapshot, IStorageDriver, IStorageDriverNodeAPI, + ListFeedMessagesOptions, + ListFeedMessagesResult, + ListFeedsOptions, + ListFeedsResult, } from "./IStorageDriver"; export type { LeasedSession } from "~/types"; diff --git a/packages/liveblocks-server/src/plugins/InMemoryDriver.ts b/packages/liveblocks-server/src/plugins/InMemoryDriver.ts index 0835bafd14f..83d623e7ac9 100644 --- a/packages/liveblocks-server/src/plugins/InMemoryDriver.ts +++ b/packages/liveblocks-server/src/plugins/InMemoryDriver.ts @@ -38,11 +38,15 @@ import type { IStorageDriver, IStorageDriverNodeAPI, LeasedSession, + ListFeedMessagesOptions, + ListFeedMessagesResult, + ListFeedsOptions, + ListFeedsResult, } from "~/interfaces"; import { NestedMap } from "~/lib/NestedMap"; import { quote } from "~/lib/text"; import { makeInMemorySnapshot } from "~/makeInMemorySnapshot"; -import type { Pos } from "~/types"; +import type { Feed, FeedMessage, Pos } from "~/types"; function buildRevNodes(nodeStream: NodeStream) { const result = new NestedMap(); @@ -126,6 +130,8 @@ export class InMemoryDriver implements IStorageDriver { private _metadb: Map; private _ydb: Map; private _leasedSessions: Map; + private _feeds: Map; + private _feedMessages: Map; // Key: `${feedId}:${messageId}` constructor(options?: { initialActor?: number; @@ -135,6 +141,8 @@ export class InMemoryDriver implements IStorageDriver { this._metadb = new Map(); this._ydb = new Map(); this._leasedSessions = new Map(); + this._feeds = new Map(); + this._feedMessages = new Map(); this._nextActor = options?.initialActor ?? -1; @@ -185,6 +193,257 @@ export class InMemoryDriver implements IStorageDriver { return 0; } + // --------------------------------------------------------------------------- + // Feed APIs + // --------------------------------------------------------------------------- + + async list_feeds(options?: ListFeedsOptions): Promise { + const limit = Math.min(options?.limit ?? 20, 100); + const since = options?.since; + const cursor = options?.cursor; + const metadata = options?.metadata; + + let feeds: Feed[] = []; + + // Collect all feeds + for (const [_, feed] of this._feeds.entries()) { + // Apply metadata filtering + if (metadata !== undefined) { + const feedMetadata = feed.metadata as JsonObject; + let matches = true; + for (const [key, value] of Object.entries(metadata)) { + if (feedMetadata[key] !== value) { + matches = false; + break; + } + } + if (!matches) continue; + } + + // Apply since filter (uses createdAt for stable ordering) + if (since !== undefined && feed.createdAt < since) { + continue; + } + + feeds.push(feed); + } + + // Sort by createdAt descending, then by feedId descending + feeds.sort((a, b) => { + if (b.createdAt !== a.createdAt) { + return b.createdAt - a.createdAt; + } + return b.feedId.localeCompare(a.feedId); + }); + + // Apply cursor-based pagination + if (cursor !== undefined) { + try { + // Simple base64url decode (using btoa/atob for browser compatibility) + const decoded = JSON.parse( + decodeURIComponent( + escape(atob(cursor.replace(/-/g, "+").replace(/_/g, "/"))) + ) + ) as [string, number]; + const [cursorFeedId, cursorCreatedAt] = decoded; + feeds = feeds.filter((s) => { + return ( + s.createdAt < cursorCreatedAt || + (s.createdAt === cursorCreatedAt && s.feedId < cursorFeedId) + ); + }); + } catch { + // Invalid cursor, ignore it + } + } + + // Apply limit (get one extra to determine if there's a next page) + let nextCursor: string | undefined; + if (feeds.length > limit) { + feeds = feeds.slice(0, limit); + // Create cursor from last returned item + const last = feeds[feeds.length - 1]; + if (last) { + const cursorData: [string, number] = [last.feedId, last.createdAt]; + // Simple base64url encode + nextCursor = btoa( + unescape(encodeURIComponent(JSON.stringify(cursorData))) + ) + .replace("/", "_") + .replace("+", "-") + .replace(/=+$/, ""); + } + } + + return { feeds, nextCursor }; + } + + async get_feed(feedId: string): Promise { + const feed = this._feeds.get(feedId); + if (feed === undefined) { + return undefined; + } + + return feed; + } + + async create_feed(feed: Feed): Promise { + // Check if feed already exists + if (this._feeds.has(feed.feedId)) { + throw new Error(`Feed ${feed.feedId} already exists`); + } + + // Store feed metadata + this._feeds.set(feed.feedId, { + feedId: feed.feedId, + metadata: feed.metadata, + createdAt: feed.createdAt, + updatedAt: feed.updatedAt, + }); + } + + async update_feed_metadata(feedId: string, metadata: Json): Promise { + const existing = this._feeds.get(feedId); + if (existing === undefined) { + throw new Error(`Feed ${feedId} not found`); + } + + this._feeds.set(feedId, { + ...existing, + metadata, + }); + } + + async delete_feed(feedId: string): Promise { + // Delete all messages for this feed + const messageKeys: string[] = []; + for (const [key] of this._feedMessages.entries()) { + if (key.startsWith(`${feedId}:`)) { + messageKeys.push(key); + } + } + for (const key of messageKeys) { + this._feedMessages.delete(key); + } + + // Delete feed + this._feeds.delete(feedId); + } + + async list_feed_messages( + feedId: string, + options?: ListFeedMessagesOptions + ): Promise { + const limit = Math.min(options?.limit ?? 20, 100); + const since = options?.since; + const cursor = options?.cursor; + + let messages: FeedMessage[] = []; + + // Collect all messages for this feed + const prefix = `${feedId}:`; + for (const [key, message] of this._feedMessages.entries()) { + if (key.startsWith(prefix)) { + // Apply since filter (uses createdAt for stable ordering) + if (since !== undefined && message.createdAt < since) { + continue; + } + + messages.push(message); + } + } + + // Sort by createdAt descending, then by id descending + messages.sort((a, b) => { + if (b.createdAt !== a.createdAt) { + return b.createdAt - a.createdAt; + } + return b.id.localeCompare(a.id); + }); + + // Apply cursor-based pagination + if (cursor !== undefined) { + try { + const decoded = JSON.parse( + decodeURIComponent( + escape(atob(cursor.replace(/-/g, "+").replace(/_/g, "/"))) + ) + ) as [string, number]; + const [cursorMessageId, cursorCreatedAt] = decoded; + messages = messages.filter((m) => { + return ( + m.createdAt < cursorCreatedAt || + (m.createdAt === cursorCreatedAt && m.id < cursorMessageId) + ); + }); + } catch { + // Invalid cursor, ignore it + } + } + + // Apply limit (get one extra to determine if there's a next page) + let nextCursor: string | undefined; + if (messages.length > limit) { + messages = messages.slice(0, limit); + // Create cursor from last returned item + const last = messages[messages.length - 1]; + if (last) { + const cursorData: [string, number] = [last.id, last.createdAt]; + nextCursor = btoa( + unescape(encodeURIComponent(JSON.stringify(cursorData))) + ) + .replace("/", "_") + .replace("+", "-") + .replace(/=+$/, ""); + } + } + + return { messages, nextCursor }; + } + + async add_feed_message(feedId: string, message: FeedMessage): Promise { + // Verify feed exists + const feed = this._feeds.get(feedId); + if (feed === undefined) { + throw new Error(`Feed ${feedId} not found`); + } + + this._feedMessages.set(`${feedId}:${message.id}`, message); + } + + async update_feed_message( + feedId: string, + messageId: string, + data: Json, + timestamp?: number + ): Promise { + const key = `${feedId}:${messageId}`; + const message = this._feedMessages.get(key); + + if (message === undefined) { + throw new Error(`Feed message ${messageId} not found in feed ${feedId}`); + } + + const effectiveTimestamp = timestamp ?? Date.now(); + if (effectiveTimestamp < message.updatedAt) { + return message; + } + + const updatedMessage: FeedMessage = { + ...message, + updatedAt: effectiveTimestamp, + data, + }; + + this._feedMessages.set(key, updatedMessage); + + return updatedMessage; + } + + async delete_feed_message(feedId: string, messageId: string): Promise { + this._feedMessages.delete(`${feedId}:${messageId}`); + } + next_actor() { return ++this._nextActor; } diff --git a/packages/liveblocks-server/src/protocol/feedErrors.ts b/packages/liveblocks-server/src/protocol/feedErrors.ts new file mode 100644 index 00000000000..c8cbee21543 --- /dev/null +++ b/packages/liveblocks-server/src/protocol/feedErrors.ts @@ -0,0 +1,70 @@ +/** + * Copyright (c) Liveblocks Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import { + FeedMsgCode, + FeedRequestErrorCode, + type FeedRequestFailedServerMsg, +} from "./feedMessages"; + +export function feedRequestFailed( + requestId: string | undefined, + code: string, + reason?: string +): FeedRequestFailedServerMsg { + return { + type: FeedMsgCode.FEED_REQUEST_FAILED, + requestId: requestId ?? "", + code, + reason, + }; +} + +/** + * Maps driver / Room errors to stable {@link FeedRequestFailedServerMsg} codes. + * Prefer matching known {@link Error#message} shapes storage drivers. + */ +export function mapFeedError(err: unknown): { code: string; reason?: string } { + if (!(err instanceof Error)) { + return { code: FeedRequestErrorCode.INTERNAL }; + } + + const m = err.message; + + if (m.includes("already exists")) { + return { code: FeedRequestErrorCode.FEED_ALREADY_EXISTS, reason: m }; + } + + if (m.includes("Feed message") && m.includes("not found")) { + return { code: FeedRequestErrorCode.FEED_MESSAGE_NOT_FOUND, reason: m }; + } + + if (m.includes("not found")) { + return { code: FeedRequestErrorCode.FEED_NOT_FOUND, reason: m }; + } + + return { code: FeedRequestErrorCode.INTERNAL, reason: m }; +} + +/** Maps an arbitrary thrown value to a {@link FeedRequestFailedServerMsg}. */ +export function feedFailureServerMsg( + requestId: string | undefined, + err: unknown +): FeedRequestFailedServerMsg { + const mapped = mapFeedError(err); + return feedRequestFailed(requestId, mapped.code, mapped.reason); +} diff --git a/packages/liveblocks-server/src/protocol/feedMessages.ts b/packages/liveblocks-server/src/protocol/feedMessages.ts new file mode 100644 index 00000000000..7ca9c36bd8a --- /dev/null +++ b/packages/liveblocks-server/src/protocol/feedMessages.ts @@ -0,0 +1,193 @@ +/** + * Copyright (c) Liveblocks Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +/** + * NOTE: this will be moved to @liveblocks/core in the future. + * Feed WebSocket message codes. + * 50x = server → client, 51x = client → server. + */ +export const FeedMsgCode = { + // Server → client (50x) + FEEDS_LIST: 500, + FEEDS_ADDED: 501, + FEEDS_UPDATED: 502, + FEED_DELETED: 503, + FEED_MESSAGES_LIST: 504, + FEED_MESSAGES_ADDED: 505, + FEED_MESSAGES_UPDATED: 506, + FEED_MESSAGES_DELETED: 507, + FEED_REQUEST_FAILED: 508, + // Client → server (51x) + FETCH_FEEDS: 510, + FETCH_FEED_MESSAGES: 511, + ADD_FEED: 512, + UPDATE_FEED: 513, + DELETE_FEED: 514, + ADD_FEED_MESSAGE: 515, + UPDATE_FEED_MESSAGE: 516, + DELETE_FEED_MESSAGE: 517, +} as const; + +/** Error codes for {@link FeedRequestFailedServerMsg}. */ +export const FeedRequestErrorCode = { + INTERNAL: "INTERNAL", + FEED_ALREADY_EXISTS: "FEED_ALREADY_EXISTS", + FEED_NOT_FOUND: "FEED_NOT_FOUND", + FEED_MESSAGE_NOT_FOUND: "FEED_MESSAGE_NOT_FOUND", +} as const; + +import type { JsonObject } from "@liveblocks/core"; + +import type { Feed, FeedMessage } from "../types"; + +// ─── Server messages (50x) ─────────────────────────────────────────────────── + +export type FeedsListServerMsg = { + type: 500; + requestId: string; + feeds: Feed[]; + nextCursor?: string; +}; + +export type FeedsAddedServerMsg = { + type: 501; + feeds: Feed[]; +}; + +export type FeedsUpdatedServerMsg = { + type: 502; + feeds: Feed[]; +}; + +export type FeedDeletedServerMsg = { + type: 503; + feedId: string; +}; + +export type FeedRequestFailedServerMsg = { + type: 508; + requestId: string; + code: string; + reason?: string; +}; + +export type FeedsServerMsg = + | FeedsListServerMsg + | FeedsAddedServerMsg + | FeedsUpdatedServerMsg + | FeedDeletedServerMsg + | FeedRequestFailedServerMsg; + +export type FeedMessagesListServerMsg = { + type: 504; + requestId: string; + feedId: string; + messages: FeedMessage[]; + nextCursor?: string; +}; + +export type FeedMessagesAddedServerMsg = { + type: 505; + feedId: string; + messages: FeedMessage[]; +}; + +export type FeedMessagesUpdatedServerMsg = { + type: 506; + feedId: string; + messages: FeedMessage[]; +}; + +export type FeedMessagesDeletedServerMsg = { + type: 507; + feedId: string; + messageIds: string[]; +}; + +export type FeedMessagesServerMsg = + | FeedMessagesListServerMsg + | FeedMessagesAddedServerMsg + | FeedMessagesUpdatedServerMsg + | FeedMessagesDeletedServerMsg + | FeedRequestFailedServerMsg; + +// ─── Client messages (51x) ─────────────────────────────────────────────────── + +export type FetchFeedsClientMsg = { + type: 510; + requestId: string; + cursor?: string; + since?: number; + limit?: number; + /** Match feeds whose metadata contains these keys (same value rules as room metadata updates). */ + metadata?: Record; +}; + +export type FetchFeedMessagesClientMsg = { + type: 511; + requestId: string; + feedId: string; + cursor?: string; + since?: number; + limit?: number; +}; + +export type AddFeedClientMsg = { + type: 512; + feedId: string; + metadata?: Record; + timestamp?: number; + requestId?: string; +}; + +export type UpdateFeedClientMsg = { + type: 513; + feedId: string; + metadata: Record; + requestId?: string; +}; + +export type DeleteFeedClientMsg = { + type: 514; + feedId: string; + requestId?: string; +}; + +export type AddFeedMessageClientMsg = { + type: 515; + feedId: string; + data: JsonObject; + id?: string; + timestamp?: number; + requestId?: string; +}; + +export type UpdateFeedMessageClientMsg = { + type: 516; + feedId: string; + messageId: string; + data: JsonObject; + timestamp?: number; + requestId?: string; +}; + +export type DeleteFeedMessageClientMsg = { + type: 517; + feedId: string; + messageId: string; + requestId?: string; +}; diff --git a/packages/liveblocks-server/src/protocol/index.ts b/packages/liveblocks-server/src/protocol/index.ts index 3e556a1ea46..1d7f27752be 100644 --- a/packages/liveblocks-server/src/protocol/index.ts +++ b/packages/liveblocks-server/src/protocol/index.ts @@ -19,6 +19,8 @@ import type { DistributiveOmit } from "@liveblocks/core"; import type { DeleteCrdtOp, SetParentKeyOp } from "./vNEXT"; +export * from "./feedErrors"; +export * from "./feedMessages"; export * from "./ProtocolVersion"; export * from "./vNEXT"; // Re-exports from @liveblocks/core, with possible additions for new protocol versions diff --git a/packages/liveblocks-server/src/types.ts b/packages/liveblocks-server/src/types.ts index 3a689ce14f4..fe93c8bc1b4 100644 --- a/packages/liveblocks-server/src/types.ts +++ b/packages/liveblocks-server/src/types.ts @@ -47,3 +47,24 @@ export type LeasedSession = { ttl: number; // time-to-live in milliseconds (default: 60000 = 1 minute) actorId: number; }; + +/** + * Feed message data structure for messages within feeds. + */ +export type FeedMessage = { + id: string; // Unique identifier + createdAt: number; // Unix timestamp in milliseconds, stable for ordering + updatedAt: number; // Unix timestamp in milliseconds, used for stale-update protection + data: Json; // Arbitrary JSON data +}; + +/** + * Feed data structure for feed-related data within a room. + * Note: Messages are stored separately and accessed via list_feed_messages. + */ +export type Feed = { + feedId: string; // Unique identifier for the feed + metadata: Json; // Arbitrary JSON metadata + createdAt: number; // Unix timestamp in milliseconds, stable for ordering + updatedAt: number; // Unix timestamp in milliseconds, same as createdAt on insert +}; diff --git a/packages/liveblocks-server/test/plugins/_generateFullTestSuite.ts b/packages/liveblocks-server/test/plugins/_generateFullTestSuite.ts index 6af8d9fe3d0..a75c1796b08 100644 --- a/packages/liveblocks-server/test/plugins/_generateFullTestSuite.ts +++ b/packages/liveblocks-server/test/plugins/_generateFullTestSuite.ts @@ -90,7 +90,7 @@ import type { UpdateObjectOp, } from "~/protocol"; import { Storage } from "~/Storage"; -import type { LeasedSession } from "~/types"; +import type { Feed, FeedMessage, LeasedSession } from "~/types"; import { YjsStorage } from "~/YjsStorage"; // Project-specific way to turn b64 string into Uint8Array @@ -684,6 +684,31 @@ export function generateArbitraries(config?: { leasedSessionPair: () => fc.tuple(arb.sessionId(), arb.leasedSession()), + feedMessage: () => { + const ts = fc.integer({ min: 0, max: Date.now() }); + return fc.record({ + id: fc.uuid(), + createdAt: ts, + updatedAt: ts, + data: arb.json(), + }); + }, + + feed: () => { + const ts = fc.integer({ min: Date.now() - 1000, max: Date.now() }); + return fc.record({ + feedId: fc.uuid(), + metadata: arb.json(), + createdAt: ts, + updatedAt: ts, + }); + }, + + feedId: () => + fc.oneof({ withCrossShrink: true }, fc.constant("feed-1"), fc.uuid()), + + feedPair: () => arb.feed().map((feed) => [feed.feedId, feed] as const), + // ------------------------------------------------------------------------- // Storage-specific arbitraries (ported from test/storage/arbitraries.ts) // ------------------------------------------------------------------------- @@ -3314,6 +3339,551 @@ export function generateFullTestSuite(config: { )); }); + describe("feed API impl", () => { + test("list_feeds on empty store is empty", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const result = await driver.list_feeds(); + expect(result.feeds).toEqual([]); + expect(result.nextCursor).toBeUndefined(); + })); + + test("get_feed on empty store is undefined", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feedId(), async (feedId) => { + expect(await driver.get_feed(feedId)).toEqual(undefined); + }) + ); + })); + + test("create_feed + get_feed", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + await driver.create_feed(feed); + expect(await driver.get_feed(feed.feedId)).toEqual(feed); + + // Cleanup + await driver.delete_feed(feed.feedId); + }) + ); + })); + + test("update_feed_metadata", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.json(), + async (feed, newMetadata) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + await driver.update_feed_metadata(feed.feedId, newMetadata); + + const updated = await driver.get_feed(feed.feedId); + expect(updated?.metadata).toEqual(newMetadata); + expect(updated?.feedId).toEqual(feed.feedId); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + expect(await driver.get_feed(feed.feedId)).toEqual(feed); + + await driver.delete_feed(feed.feedId); + expect(await driver.get_feed(feed.feedId)).toBeUndefined(); + }) + ); + })); + + test("delete_feed on non-existent feed is no-op", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feedId(), async (feedId) => { + await driver.delete_feed(feedId); + expect(await driver.get_feed(feedId)).toBeUndefined(); + }) + ); + })); + + test("list_feeds returns all feeds", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + fc + .array(arb.feedPair()) + .map((x) => new Map(x)) + .filter((m) => m.size > 0 && m.size <= 5), + + async (entries) => { + // Ensure feeds don't already exist + for (const [feedId] of entries) { + await driver.delete_feed(feedId); + } + + // Create all feeds + for (const [, feed] of entries) { + await driver.create_feed(feed); + } + + // List all feeds + const result = await driver.list_feeds(); + + // Should return all feeds + expect(result.feeds.length).toEqual(entries.size); + const feedIds = new Set(result.feeds.map((f) => f.feedId)); + for (const [feedId] of entries) { + expect(feedIds.has(feedId)).toBe(true); + } + + // Cleanup + for (const [feedId] of entries) { + await driver.delete_feed(feedId); + } + } + ) + ); + })); + + test("list_feeds with pagination", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + // Create 3 feeds with different createdAt timestamps + const feeds: Feed[] = [ + { + feedId: "feed-1", + metadata: {}, + createdAt: 1000, + updatedAt: 1000, + }, + { + feedId: "feed-2", + metadata: {}, + createdAt: 2000, + updatedAt: 2000, + }, + { + feedId: "feed-3", + metadata: {}, + createdAt: 3000, + updatedAt: 3000, + }, + ]; + + for (const feed of feeds) { + await driver.create_feed(feed); + } + + // List with limit + const page1 = await driver.list_feeds({ limit: 2 }); + expect(page1.feeds.length).toBe(2); + expect(page1.nextCursor).toBeDefined(); + + // Should be sorted by createdAt descending (newest first) + expect(page1.feeds[0]?.feedId).toBe("feed-3"); + expect(page1.feeds[1]?.feedId).toBe("feed-2"); + + // Get next page + const page2 = await driver.list_feeds({ + cursor: page1.nextCursor, + }); + expect(page2.feeds.length).toBe(1); + expect(page2.feeds[0]?.feedId).toBe("feed-1"); + expect(page2.nextCursor).toBeUndefined(); + + // Cleanup + for (const feed of feeds) { + await driver.delete_feed(feed.feedId); + } + })); + + test("list_feeds with since filter", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feeds: Feed[] = [ + { + feedId: "feed-old", + metadata: {}, + createdAt: 1000, + updatedAt: 1000, + }, + { + feedId: "feed-new", + metadata: {}, + createdAt: 5000, + updatedAt: 5000, + }, + ]; + + for (const feed of feeds) { + await driver.create_feed(feed); + } + + // List feeds since 2000 + const result = await driver.list_feeds({ since: 2000 }); + expect(result.feeds.length).toBe(1); + expect(result.feeds[0]?.feedId).toBe("feed-new"); + + // Cleanup + for (const feed of feeds) { + await driver.delete_feed(feed.feedId); + } + })); + + test("add_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + // Create feed + await driver.create_feed(feed); + + // Add message + await driver.add_feed_message(feed.feedId, message); + + // List messages and verify message was added + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages).toHaveLength(1); + expect(result.messages[0]).toEqual(message); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("list_feed_messages", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(0); + + // Cleanup + await driver.delete_feed(feed.feedId); + }) + ); + })); + + test("list_feed_messages with pagination", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const messages: FeedMessage[] = [ + { + id: "msg-1", + createdAt: 1000, + updatedAt: 1000, + data: { value: 1 }, + }, + { + id: "msg-2", + createdAt: 2000, + updatedAt: 2000, + data: { value: 2 }, + }, + { + id: "msg-3", + createdAt: 3000, + updatedAt: 3000, + data: { value: 3 }, + }, + ]; + + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + + await driver.create_feed(feed); + + // Add messages one by one + for (const message of messages) { + await driver.add_feed_message(feed.feedId, message); + } + + // List with limit + const page1 = await driver.list_feed_messages(feed.feedId, { + limit: 2, + }); + expect(page1.messages.length).toBe(2); + expect(page1.nextCursor).toBeDefined(); + + // Should be sorted by createdAt descending + expect(page1.messages[0]?.id).toBe("msg-3"); + expect(page1.messages[1]?.id).toBe("msg-2"); + + // Get next page + const page2 = await driver.list_feed_messages(feed.feedId, { + cursor: page1.nextCursor, + }); + expect(page2.messages.length).toBe(1); + expect(page2.messages[0]?.id).toBe("msg-1"); + expect(page2.nextCursor).toBeUndefined(); + + // Cleanup + await driver.delete_feed(feed.feedId); + })); + + test("update_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + arb.json(), + async (feed, message, newData) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.update_feed_message( + feed.feedId, + message.id, + newData + ); + + const result = await driver.list_feed_messages(feed.feedId); + const updatedMessage = result.messages.find( + (m) => m.id === message.id + ); + expect(updatedMessage?.data).toEqual(newData); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("update_feed_message ignores stale timestamped updates for same message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + const message: FeedMessage = { + id: "msg-1", + createdAt: 1000, + updatedAt: 1000, + data: { value: "original" }, + }; + + await driver.delete_feed(feed.feedId); + await driver.create_feed(feed); + await driver.add_feed_message(feed.feedId, message); + + const newer = await driver.update_feed_message( + feed.feedId, + message.id, + { value: "newer-update" }, + 2000 + ); + expect(newer.updatedAt).toBe(2000); + expect(newer.data).toEqual({ value: "newer-update" }); + + const stale = await driver.update_feed_message( + feed.feedId, + message.id, + { value: "stale-update" }, + 1500 + ); + expect(stale.updatedAt).toBe(2000); + expect(stale.data).toEqual({ value: "newer-update" }); + + const result = await driver.list_feed_messages(feed.feedId); + const latest = result.messages.find((m) => m.id === message.id); + expect(latest?.updatedAt).toBe(2000); + expect(latest?.data).toEqual({ value: "newer-update" }); + + await driver.delete_feed(feed.feedId); + })); + + test("list_feed_messages order remains by createdAt after update", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + const messages: FeedMessage[] = [ + { + id: "msg-first", + createdAt: 1000, + updatedAt: 1000, + data: { value: "first" }, + }, + { + id: "msg-second", + createdAt: 2000, + updatedAt: 2000, + data: { value: "second" }, + }, + ]; + + await driver.delete_feed(feed.feedId); + await driver.create_feed(feed); + for (const m of messages) { + await driver.add_feed_message(feed.feedId, m); + } + + const beforeUpdate = await driver.list_feed_messages(feed.feedId); + expect(beforeUpdate.messages[0]?.id).toBe("msg-second"); + expect(beforeUpdate.messages[1]?.id).toBe("msg-first"); + + await driver.update_feed_message( + feed.feedId, + "msg-first", + { value: "first-updated" }, + 9999 + ); + + const afterUpdate = await driver.list_feed_messages(feed.feedId); + expect(afterUpdate.messages[0]?.id).toBe("msg-second"); + expect(afterUpdate.messages[1]?.id).toBe("msg-first"); + expect(afterUpdate.messages[1]?.data).toEqual({ + value: "first-updated", + }); + expect(afterUpdate.messages[1]?.updatedAt).toBe(9999); + + await driver.delete_feed(feed.feedId); + })); + + test("delete_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed_message(feed.feedId, message.id); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(0); + expect( + result.messages.find((m) => m.id === message.id) + ).toBeUndefined(); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed_message on non-existent message is no-op", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed_message( + feed.feedId, + "non-existent-message-id" + ); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(1); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed deletes all messages", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed(feed.feedId); + + // Feed should be gone + expect(await driver.get_feed(feed.feedId)).toBeUndefined(); + + // Messages should be gone too + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages).toEqual([]); + } + ) + ); + })); + }); + describe("Storage behavior (high-level)", () => { /** * Helper to create a Storage instance backed by the driver. The raw nodes diff --git a/tools/liveblocks-cli/src/dev-server/db/BunSQLiteDriver.ts b/tools/liveblocks-cli/src/dev-server/db/BunSQLiteDriver.ts index 667a5e098fa..4dfd137c5c9 100644 --- a/tools/liveblocks-cli/src/dev-server/db/BunSQLiteDriver.ts +++ b/tools/liveblocks-cli/src/dev-server/db/BunSQLiteDriver.ts @@ -29,10 +29,16 @@ import type { } from "@liveblocks/core"; import { asPos, CrdtType, nn } from "@liveblocks/core"; import type { + Feed, + FeedMessage, IReadableSnapshot, IStorageDriver, IStorageDriverNodeAPI, LeasedSession, + ListFeedMessagesOptions, + ListFeedMessagesResult, + ListFeedsOptions, + ListFeedsResult, Logger, Pos, YDocId, @@ -43,7 +49,7 @@ import { plainLsonToNodeStream, quote, } from "@liveblocks/server"; -import { Database } from "bun:sqlite"; +import { Database, type SQLQueryBindings } from "bun:sqlite"; function tryParseJson( value: string | undefined @@ -76,6 +82,21 @@ type YdocsRow = { data: Uint8Array; }; +type FeedRow = { + feed_id: string; + jmetadata: string; + created_at: number; + updated_at: number; +}; + +type FeedMessageRow = { + feed_id: string; + message_id: string; + jdata: string; + created_at: number; + updated_at: number; +}; + type LeasedSessionRow = { session_id: string; jpresence: string; // JSON @@ -317,6 +338,7 @@ export class BunSQLiteDriver implements IStorageDriver { db.run("PRAGMA journal_mode = WAL"); db.run("PRAGMA case_sensitive_like = ON"); + db.run("PRAGMA foreign_keys = ON"); // Create a table for room info db.run( @@ -370,6 +392,35 @@ export class BunSQLiteDriver implements IStorageDriver { "CREATE INDEX IF NOT EXISTS idx_leased_sessions_expiry ON leased_sessions(updated_at, ttl)" ); + // Create a table for feeds + db.run( + `CREATE TABLE IF NOT EXISTS feeds ( + feed_id TEXT NOT NULL PRIMARY KEY, + jmetadata TEXT NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + )` + ); + db.run( + "CREATE INDEX IF NOT EXISTS idx_feeds_created_at ON feeds(created_at DESC, feed_id DESC)" + ); + + // Create a table for feed messages + db.run( + `CREATE TABLE IF NOT EXISTS feed_messages ( + feed_id TEXT NOT NULL, + message_id TEXT NOT NULL, + jdata TEXT NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL, + PRIMARY KEY (feed_id, message_id), + FOREIGN KEY (feed_id) REFERENCES feeds(feed_id) ON DELETE CASCADE + )` + ); + db.run( + "CREATE INDEX IF NOT EXISTS idx_feed_messages_feed_created ON feed_messages(feed_id, created_at DESC, message_id DESC)" + ); + this.db = db; } @@ -807,6 +858,286 @@ export class BunSQLiteDriver implements IStorageDriver { .run(sessionId); } + // --------------------------------------------------------------------------- + // Feed APIs + // --------------------------------------------------------------------------- + + list_feeds(options?: ListFeedsOptions): ListFeedsResult { + const limit = Math.min(options?.limit ?? 20, 100); + const since = options?.since; + const cursor = options?.cursor; + const metadata = options?.metadata; + + let query = + "SELECT feed_id, jmetadata, created_at, updated_at FROM feeds WHERE 1=1"; + const params: SQLQueryBindings[] = []; + + if (metadata !== undefined) { + for (const [key, value] of Object.entries(metadata)) { + const sqlValue = typeof value === "boolean" ? Number(value) : value; + query += " AND json_extract(jmetadata, '$.' || ?) = ?"; + params.push(key, sqlValue as SQLQueryBindings); + } + } + + if (since !== undefined) { + query += " AND created_at >= ?"; + params.push(since); + } + + if (cursor !== undefined) { + try { + const decoded = JSON.parse( + Buffer.from( + cursor.replace(/-/g, "+").replace(/_/g, "/"), + "base64" + ).toString("utf8") + ) as [string, number]; + const [feedId, createdAt] = decoded; + query += " AND (created_at < ? OR (created_at = ? AND feed_id < ?))"; + params.push(createdAt, createdAt, feedId); + } catch { + // Invalid cursor, ignore it + } + } + + query += " ORDER BY created_at DESC, feed_id DESC LIMIT ?"; + params.push(limit + 1); + + const rows = this.db + .query(query) + .all(...params); + + let nextCursor: string | undefined; + if (rows.length > limit) { + rows.pop(); + const last = rows[rows.length - 1]; + if (last) { + const cursorData: [string, number] = [last.feed_id, last.created_at]; + nextCursor = Buffer.from(JSON.stringify(cursorData), "utf8") + .toString("base64") + .replace(/\+/g, "-") + .replace(/\//g, "_") + .replace(/=+$/, ""); + } + } + + const feeds: Feed[] = rows.map((row) => ({ + feedId: row.feed_id, + metadata: JSON.parse(row.jmetadata) as Feed["metadata"], + createdAt: row.created_at, + updatedAt: row.updated_at, + })); + + return { feeds, nextCursor }; + } + + get_feed(feedId: string): Feed | undefined { + const row = this.db + .query( + "SELECT feed_id, jmetadata, created_at, updated_at FROM feeds WHERE feed_id = ?" + ) + .get(feedId); + if (row === undefined || row === null) return undefined; + return { + feedId: row.feed_id, + metadata: JSON.parse(row.jmetadata) as Feed["metadata"], + createdAt: row.created_at, + updatedAt: row.updated_at, + }; + } + + create_feed(feed: Feed): void { + const existing = this.db + .query, [string]>( + "SELECT feed_id FROM feeds WHERE feed_id = ?" + ) + .get(feed.feedId); + if (existing !== undefined && existing !== null) { + throw new Error(`Feed ${feed.feedId} already exists`); + } + this.db + .query( + "INSERT INTO feeds (feed_id, jmetadata, created_at, updated_at) VALUES (?, ?, ?, ?)" + ) + .run( + feed.feedId, + JSON.stringify(feed.metadata), + feed.createdAt, + feed.updatedAt + ); + } + + update_feed_metadata(feedId: string, metadata: Feed["metadata"]): void { + const result = this.db + .query( + "UPDATE feeds SET jmetadata = ? WHERE feed_id = ? RETURNING feed_id, jmetadata, created_at, updated_at" + ) + .get(JSON.stringify(metadata), feedId); + if (result === undefined || result === null) { + throw new Error(`Feed ${feedId} not found`); + } + } + + delete_feed(feedId: string): void { + this.db + .query("DELETE FROM feeds WHERE feed_id = ?") + .run(feedId); + } + + list_feed_messages( + feedId: string, + options?: ListFeedMessagesOptions + ): ListFeedMessagesResult { + const limit = Math.min(options?.limit ?? 20, 100); + const since = options?.since; + const cursor = options?.cursor; + + let query = + "SELECT feed_id, message_id, jdata, created_at, updated_at FROM feed_messages WHERE feed_id = ?"; + const params: SQLQueryBindings[] = [feedId]; + + if (since !== undefined) { + query += " AND created_at >= ?"; + params.push(since); + } + + if (cursor !== undefined) { + try { + const decoded = JSON.parse( + Buffer.from( + cursor.replace(/-/g, "+").replace(/_/g, "/"), + "base64" + ).toString("utf8") + ) as [string, number]; + const [messageId, createdAt] = decoded; + query += " AND (created_at < ? OR (created_at = ? AND message_id < ?))"; + params.push(createdAt, createdAt, messageId); + } catch { + // Invalid cursor, ignore it + } + } + + query += " ORDER BY created_at DESC, message_id DESC LIMIT ?"; + params.push(limit + 1); + + const rows = this.db + .query(query) + .all(...params); + + let nextCursor: string | undefined; + if (rows.length > limit) { + rows.pop(); + const last = rows[rows.length - 1]; + if (last) { + const cursorData: [string, number] = [last.message_id, last.created_at]; + nextCursor = Buffer.from(JSON.stringify(cursorData), "utf8") + .toString("base64") + .replace(/\+/g, "-") + .replace(/\//g, "_") + .replace(/=+$/, ""); + } + } + + const messages: FeedMessage[] = rows.map((row) => ({ + id: row.message_id, + data: JSON.parse(row.jdata) as FeedMessage["data"], + createdAt: row.created_at, + updatedAt: row.updated_at, + })); + + return { messages, nextCursor }; + } + + add_feed_message(feedId: string, message: FeedMessage): void { + const feed = this.get_feed(feedId); + if (feed === undefined) { + throw new Error(`Feed ${feedId} not found`); + } + this.db + .query( + "INSERT INTO feed_messages (feed_id, message_id, jdata, created_at, updated_at) VALUES (?, ?, ?, ?, ?)" + ) + .run( + feedId, + message.id, + JSON.stringify(message.data), + message.createdAt, + message.updatedAt + ); + } + + update_feed_message( + feedId: string, + messageId: string, + data: FeedMessage["data"], + timestamp?: number + ): FeedMessage { + const existing = this.db + .query( + "SELECT feed_id, message_id, jdata, created_at, updated_at FROM feed_messages WHERE feed_id = ? AND message_id = ?" + ) + .get(feedId, messageId); + if (existing === undefined || existing === null) { + throw new Error(`Feed message ${messageId} not found in feed ${feedId}`); + } + + const effectiveTimestamp = timestamp ?? Date.now(); + if (effectiveTimestamp < existing.updated_at) { + return { + id: existing.message_id, + data: JSON.parse(existing.jdata) as FeedMessage["data"], + createdAt: existing.created_at, + updatedAt: existing.updated_at, + }; + } + + const result = this.db + .query< + FeedMessageRow, + [string, number, string, string, number] + >( + "UPDATE feed_messages SET jdata = ?, updated_at = ? WHERE feed_id = ? AND message_id = ? AND updated_at <= ? RETURNING feed_id, message_id, jdata, created_at, updated_at" + ) + .get( + JSON.stringify(data), + effectiveTimestamp, + feedId, + messageId, + effectiveTimestamp + ); + if (result === undefined || result === null) { + const latest = this.db + .query( + "SELECT feed_id, message_id, jdata, created_at, updated_at FROM feed_messages WHERE feed_id = ? AND message_id = ?" + ) + .get(feedId, messageId); + if (latest === undefined || latest === null) { + throw new Error(`Feed message ${messageId} not found in feed ${feedId}`); + } + return { + id: latest.message_id, + data: JSON.parse(latest.jdata) as FeedMessage["data"], + createdAt: latest.created_at, + updatedAt: latest.updated_at, + }; + } + return { + id: result.message_id, + data: JSON.parse(result.jdata) as FeedMessage["data"], + createdAt: result.created_at, + updatedAt: result.updated_at, + }; + } + + delete_feed_message(feedId: string, messageId: string): void { + this.db + .query( + "DELETE FROM feed_messages WHERE feed_id = ? AND message_id = ?" + ) + .run(feedId, messageId); + } + close() { this.db.close(); } diff --git a/tools/liveblocks-cli/src/dev-server/routes/rest-api.ts b/tools/liveblocks-cli/src/dev-server/routes/rest-api.ts index 560aa349e7b..4e3bd80ebb5 100644 --- a/tools/liveblocks-cli/src/dev-server/routes/rest-api.ts +++ b/tools/liveblocks-cli/src/dev-server/routes/rest-api.ts @@ -43,7 +43,7 @@ import type { DbRoom, RoomFilters } from "~/dev-server/db/rooms"; import * as Rooms from "~/dev-server/db/rooms"; import { authorizeSecretKey } from "~/dev-server/lib/auth"; import { yDocToJson } from "~/dev-server/lib/ydoc"; -import { NOT_IMPLEMENTED } from "~/dev-server/responses"; +import { DUMMY, NOT_IMPLEMENTED } from "~/dev-server/responses"; enum SerializationFormat { PlainLson = "plain-lson", // the default @@ -476,6 +476,117 @@ zen.route("GET /v2/rooms//active_users", ({ p }) => { return json({ data }); }); +zen.route("DELETE /v2/rooms//feeds/", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + ok: true, + }); +}); +zen.route( + "DELETE /v2/rooms//feeds//messages/", + ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + ok: true, + }); + } +); +zen.route("GET /v2/rooms//feeds", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + feeds: [ + { + feedId: "123", + metadata: { + title: "Test Feed", + description: "This is a test feed", + }, + timestamp: new Date().getTime(), + }, + ], + nextCursor: null, + }); +}); +zen.route("GET /v2/rooms//feeds/", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + feedId: "123", + timestamp: new Date().getTime(), + metadata: { + title: "Test Feed", + description: "This is a test feed", + }, + }); +}); +zen.route("GET /v2/rooms//feeds//messages", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + messages: [ + { + messageId: "123", + data: { + content: "This is a test message", + }, + }, + ], + nextCursor: null, + }); +}); +zen.route("PATCH /v2/rooms//feeds/", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + ok: true, + }); +}); +zen.route( + "PATCH /v2/rooms//feeds//messages/", + ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + ok: true, + }); + } +); +zen.route("POST /v2/rooms//feed", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + feedId: "123", + timestamp: new Date().getTime(), + metadata: { + title: "Test Feed", + description: "This is a test feed", + }, + }); +}); +zen.route("POST /v2/rooms//feeds//messages", ({ p }) => { + if (!Rooms.getRoom(p.roomId)) { + throw ROOM_NOT_FOUND(p.roomId); + } + return DUMMY({ + id: "123", + timestamp: new Date().getTime(), + data: { + content: "This is a test message", + }, + }); +}); + /** * ------------------------------------------------------------ * NOT IMPLEMENTED ROUTES diff --git a/tools/liveblocks-cli/test/plugins/_generateFullTestSuite.ts b/tools/liveblocks-cli/test/plugins/_generateFullTestSuite.ts index 70810fcc91d..0d79ef9fefc 100644 --- a/tools/liveblocks-cli/test/plugins/_generateFullTestSuite.ts +++ b/tools/liveblocks-cli/test/plugins/_generateFullTestSuite.ts @@ -73,6 +73,8 @@ import type { CreateRegisterOp, DeleteCrdtOp, DeleteObjectKeyOp, + Feed, + FeedMessage, Guid, HasOpId, IStorageDriver, @@ -689,6 +691,31 @@ export function generateArbitraries(config?: { leasedSessionPair: () => fc.tuple(arb.sessionId(), arb.leasedSession()), + feedMessage: () => { + const ts = fc.integer({ min: 0, max: Date.now() }); + return fc.record({ + id: fc.uuid(), + createdAt: ts, + updatedAt: ts, + data: arb.json(), + }); + }, + + feed: () => { + const ts = fc.integer({ min: Date.now() - 1000, max: Date.now() }); + return fc.record({ + feedId: fc.uuid(), + metadata: arb.json(), + createdAt: ts, + updatedAt: ts, + }); + }, + + feedId: () => + fc.oneof({ withCrossShrink: true }, fc.constant("feed-1"), fc.uuid()), + + feedPair: () => arb.feed().map((feed) => [feed.feedId, feed] as const), + // ------------------------------------------------------------------------- // Storage-specific arbitraries (ported from test/storage/arbitraries.ts) // ------------------------------------------------------------------------- @@ -3319,6 +3346,551 @@ export function generateFullTestSuite(config: { )); }); + describe("feed API impl", () => { + test("list_feeds on empty store is empty", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const result = await driver.list_feeds(); + expect(result.feeds).toEqual([]); + expect(result.nextCursor).toBeUndefined(); + })); + + test("get_feed on empty store is undefined", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feedId(), async (feedId) => { + expect(await driver.get_feed(feedId)).toEqual(undefined); + }) + ); + })); + + test("create_feed + get_feed", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + await driver.create_feed(feed); + expect(await driver.get_feed(feed.feedId)).toEqual(feed); + + // Cleanup + await driver.delete_feed(feed.feedId); + }) + ); + })); + + test("update_feed_metadata", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.json(), + async (feed, newMetadata) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + await driver.update_feed_metadata(feed.feedId, newMetadata); + + const updated = await driver.get_feed(feed.feedId); + expect(updated?.metadata).toEqual(newMetadata); + expect(updated?.feedId).toEqual(feed.feedId); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + expect(await driver.get_feed(feed.feedId)).toEqual(feed); + + await driver.delete_feed(feed.feedId); + expect(await driver.get_feed(feed.feedId)).toBeUndefined(); + }) + ); + })); + + test("delete_feed on non-existent feed is no-op", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feedId(), async (feedId) => { + await driver.delete_feed(feedId); + expect(await driver.get_feed(feedId)).toBeUndefined(); + }) + ); + })); + + test("list_feeds returns all feeds", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + fc + .array(arb.feedPair()) + .map((x) => new Map(x)) + .filter((m) => m.size > 0 && m.size <= 5), + + async (entries) => { + // Ensure feeds don't already exist + for (const [feedId] of entries) { + await driver.delete_feed(feedId); + } + + // Create all feeds + for (const [, feed] of entries) { + await driver.create_feed(feed); + } + + // List all feeds + const result = await driver.list_feeds(); + + // Should return all feeds + expect(result.feeds.length).toEqual(entries.size); + const feedIds = new Set(result.feeds.map((f) => f.feedId)); + for (const [feedId] of entries) { + expect(feedIds.has(feedId)).toBe(true); + } + + // Cleanup + for (const [feedId] of entries) { + await driver.delete_feed(feedId); + } + } + ) + ); + })); + + test("list_feeds with pagination", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + // Create 3 feeds with different createdAt timestamps + const feeds: Feed[] = [ + { + feedId: "feed-1", + metadata: {}, + createdAt: 1000, + updatedAt: 1000, + }, + { + feedId: "feed-2", + metadata: {}, + createdAt: 2000, + updatedAt: 2000, + }, + { + feedId: "feed-3", + metadata: {}, + createdAt: 3000, + updatedAt: 3000, + }, + ]; + + for (const feed of feeds) { + await driver.create_feed(feed); + } + + // List with limit + const page1 = await driver.list_feeds({ limit: 2 }); + expect(page1.feeds.length).toBe(2); + expect(page1.nextCursor).toBeDefined(); + + // Should be sorted by createdAt descending (newest first) + expect(page1.feeds[0]?.feedId).toBe("feed-3"); + expect(page1.feeds[1]?.feedId).toBe("feed-2"); + + // Get next page + const page2 = await driver.list_feeds({ + cursor: page1.nextCursor, + }); + expect(page2.feeds.length).toBe(1); + expect(page2.feeds[0]?.feedId).toBe("feed-1"); + expect(page2.nextCursor).toBeUndefined(); + + // Cleanup + for (const feed of feeds) { + await driver.delete_feed(feed.feedId); + } + })); + + test("list_feeds with since filter", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feeds: Feed[] = [ + { + feedId: "feed-old", + metadata: {}, + createdAt: 1000, + updatedAt: 1000, + }, + { + feedId: "feed-new", + metadata: {}, + createdAt: 5000, + updatedAt: 5000, + }, + ]; + + for (const feed of feeds) { + await driver.create_feed(feed); + } + + // List feeds since 2000 + const result = await driver.list_feeds({ since: 2000 }); + expect(result.feeds.length).toBe(1); + expect(result.feeds[0]?.feedId).toBe("feed-new"); + + // Cleanup + for (const feed of feeds) { + await driver.delete_feed(feed.feedId); + } + })); + + test("add_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + // Create feed + await driver.create_feed(feed); + + // Add message + await driver.add_feed_message(feed.feedId, message); + + // List messages and verify message was added + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages).toHaveLength(1); + expect(result.messages[0]).toEqual(message); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("list_feed_messages", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty(arb.feed(), async (feed) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(0); + + // Cleanup + await driver.delete_feed(feed.feedId); + }) + ); + })); + + test("list_feed_messages with pagination", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const messages: FeedMessage[] = [ + { + id: "msg-1", + createdAt: 1000, + updatedAt: 1000, + data: { value: 1 }, + }, + { + id: "msg-2", + createdAt: 2000, + updatedAt: 2000, + data: { value: 2 }, + }, + { + id: "msg-3", + createdAt: 3000, + updatedAt: 3000, + data: { value: 3 }, + }, + ]; + + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + + await driver.create_feed(feed); + + // Add messages one by one + for (const message of messages) { + await driver.add_feed_message(feed.feedId, message); + } + + // List with limit + const page1 = await driver.list_feed_messages(feed.feedId, { + limit: 2, + }); + expect(page1.messages.length).toBe(2); + expect(page1.nextCursor).toBeDefined(); + + // Should be sorted by createdAt descending + expect(page1.messages[0]?.id).toBe("msg-3"); + expect(page1.messages[1]?.id).toBe("msg-2"); + + // Get next page + const page2 = await driver.list_feed_messages(feed.feedId, { + cursor: page1.nextCursor, + }); + expect(page2.messages.length).toBe(1); + expect(page2.messages[0]?.id).toBe("msg-1"); + expect(page2.nextCursor).toBeUndefined(); + + // Cleanup + await driver.delete_feed(feed.feedId); + })); + + test("update_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + arb.json(), + async (feed, message, newData) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.update_feed_message( + feed.feedId, + message.id, + newData + ); + + const result = await driver.list_feed_messages(feed.feedId); + const updatedMessage = result.messages.find( + (m) => m.id === message.id + ); + expect(updatedMessage?.data).toEqual(newData); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("update_feed_message ignores stale timestamped updates for same message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + const message: FeedMessage = { + id: "msg-1", + createdAt: 1000, + updatedAt: 1000, + data: { value: "original" }, + }; + + await driver.delete_feed(feed.feedId); + await driver.create_feed(feed); + await driver.add_feed_message(feed.feedId, message); + + const newer = await driver.update_feed_message( + feed.feedId, + message.id, + { value: "newer-update" }, + 2000 + ); + expect(newer.updatedAt).toBe(2000); + expect(newer.data).toEqual({ value: "newer-update" }); + + const stale = await driver.update_feed_message( + feed.feedId, + message.id, + { value: "stale-update" }, + 1500 + ); + expect(stale.updatedAt).toBe(2000); + expect(stale.data).toEqual({ value: "newer-update" }); + + const result = await driver.list_feed_messages(feed.feedId); + const latest = result.messages.find((m) => m.id === message.id); + expect(latest?.updatedAt).toBe(2000); + expect(latest?.data).toEqual({ value: "newer-update" }); + + await driver.delete_feed(feed.feedId); + })); + + test("list_feed_messages order remains by createdAt after update", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + const feed: Feed = { + feedId: "test-feed", + metadata: {}, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + const messages: FeedMessage[] = [ + { + id: "msg-first", + createdAt: 1000, + updatedAt: 1000, + data: { value: "first" }, + }, + { + id: "msg-second", + createdAt: 2000, + updatedAt: 2000, + data: { value: "second" }, + }, + ]; + + await driver.delete_feed(feed.feedId); + await driver.create_feed(feed); + for (const m of messages) { + await driver.add_feed_message(feed.feedId, m); + } + + const beforeUpdate = await driver.list_feed_messages(feed.feedId); + expect(beforeUpdate.messages[0]?.id).toBe("msg-second"); + expect(beforeUpdate.messages[1]?.id).toBe("msg-first"); + + await driver.update_feed_message( + feed.feedId, + "msg-first", + { value: "first-updated" }, + 9999 + ); + + const afterUpdate = await driver.list_feed_messages(feed.feedId); + expect(afterUpdate.messages[0]?.id).toBe("msg-second"); + expect(afterUpdate.messages[1]?.id).toBe("msg-first"); + expect(afterUpdate.messages[1]?.data).toEqual({ + value: "first-updated", + }); + expect(afterUpdate.messages[1]?.updatedAt).toBe(9999); + + await driver.delete_feed(feed.feedId); + })); + + test("delete_feed_message", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed_message(feed.feedId, message.id); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(0); + expect( + result.messages.find((m) => m.id === message.id) + ).toBeUndefined(); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed_message on non-existent message is no-op", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + + // Add message to the feed + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed_message( + feed.feedId, + "non-existent-message-id" + ); + + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages.length).toBe(1); + + // Cleanup + await driver.delete_feed(feed.feedId); + } + ) + ); + })); + + test("delete_feed deletes all messages", () => + runTest(async (driver) => { + if (config.name === "dos-kv") return Promise.resolve(); + return fc.assert( + fc.asyncProperty( + arb.feed(), + arb.feedMessage(), + async (feed, message) => { + // Ensure feed doesn't already exist + await driver.delete_feed(feed.feedId); + + await driver.create_feed(feed); + await driver.add_feed_message(feed.feedId, message); + + await driver.delete_feed(feed.feedId); + + // Feed should be gone + expect(await driver.get_feed(feed.feedId)).toBeUndefined(); + + // Messages should be gone too + const result = await driver.list_feed_messages(feed.feedId); + expect(result.messages).toEqual([]); + } + ) + ); + })); + }); + describe("Storage behavior (high-level)", () => { /** * Helper to create a Storage instance backed by the driver. The raw nodes