Skip to content
37 changes: 23 additions & 14 deletions backend/src/lib/compression.ts
Original file line number Diff line number Diff line change
@@ -1,23 +1,23 @@
/* Copyright Contributors to the Open Cluster Management project */
import type { Readable, Transform } from 'node:stream'
import { pipeline } from 'node:stream'
import { promisify } from 'node:util'
import type { Zlib } from 'node:zlib'
import {
createBrotliCompress,
createBrotliDecompress,
createDeflate,
createGunzip,
createGzip,
createInflate,
inflateRaw,
deflateRaw,
type Zlib,
inflateRaw,
} from 'node:zlib'
import { logger } from './logger'
import type { ServerSideEvent, WatchEvent } from './server-side-events'
import { getEventDict } from '../routes/events'
import { getAppDict, type ICompressedResource, type ITransformedResource } from '../routes/aggregators/applications'
import { promisify } from 'node:util'
import { getEventDict } from '../routes/events'
import type { IResource } from './../resources/resource'
import { logger } from './logger'
import type { ServerSideEvent, WatchEvent } from './server-side-events'

const MAX_RECENTLY_ADDED = 200

Expand Down Expand Up @@ -138,7 +138,7 @@ export class FifoSet<T> {
const bigStrings: FifoSet<string> = new FifoSet(200)

export async function deflateResource(resource: IResource, dictionary: Dictionary): Promise<Buffer> {
const res = compressResource(resource as UncompressedResourceType, dictionary)
const res = compressResource(resource, dictionary)
let buffer
try {
buffer = await promisify(deflateRaw)(JSON.stringify(res))
Expand All @@ -155,7 +155,6 @@ export async function deflateResource(resource: IResource, dictionary: Dictionar
function compressResource(resource: UncompressedResourceType, dictionary: Dictionary): CompressedResourceType {
if (resource) {
if (Array.isArray(resource)) {
// eslint-disable-next-line @typescript-eslint/no-unsafe-return, @typescript-eslint/no-unsafe-call
return resource.map((item: UncompressedResourceType) => compressResource(item, dictionary))
} else if (typeof resource === 'object') {
// eslint-disable-next-line @typescript-eslint/no-explicit-any
Expand All @@ -176,8 +175,9 @@ function compressResource(resource: UncompressedResourceType, dictionary: Dictio
res[dictionary.add(key)] = resource[key]
} else {
const inx = dictionary.add(key)
if (valueInDictionaryKeys.has(key)) {
res[inx] = dictionary.add(resource[key] as string)
// Guard against non-string values (e.g. nested CRD OpenAPI schema objects) corrupting the shared dictionary.
if (valueInDictionaryKeys.has(key) && typeof resource[key] === 'string') {
res[inx] = dictionary.add(resource[key])
} else {
res[inx] = compressResource(resource[key] as UncompressedResourceType, dictionary)
}
Expand Down Expand Up @@ -230,7 +230,7 @@ function compressResource(resource: UncompressedResourceType, dictionary: Dictio
export async function inflateResource(buffer: Buffer, dictionary: Dictionary): Promise<IResource> {
let inflated
try {
inflated = (await promisify(inflateRaw)(buffer)).toString()
inflated = (await promisify(inflateRaw)(new Uint8Array(buffer))).toString()
} catch (err: unknown) {
logger.error({
msg: 'Error from inflateRaw during inflateResource',
Expand All @@ -243,11 +243,18 @@ export async function inflateResource(buffer: Buffer, dictionary: Dictionary): P
}

export async function inflateEvent(event: ServerSideEvent): Promise<ServerSideEvent> {
const { id, data } = event
const { type, object } = data as WatchEvent
const { id, name, namespace, data } = event
if (!data || typeof data !== 'object') return event
const watchEvent = data as WatchEvent & { meta?: unknown }
const { type, object } = watchEvent
return !object
? event
: { id, data: { type, object: Buffer.isBuffer(object) ? await inflateResource(object, getEventDict()) : object } }
: {
id,
name,
namespace,
data: { type, object: Buffer.isBuffer(object) ? await inflateResource(object, getEventDict()) : object },
}
}

export async function inflateApps(apps: ICompressedResource[]): Promise<ITransformedResource[]> {
Expand Down Expand Up @@ -278,6 +285,8 @@ function decompressResource(resource: CompressedResourceType, dictionary: Dictio
for (const inx in resource) {
if (Object.prototype.hasOwnProperty.call(resource, inx)) {
const key = dictionary.get(Number(inx))
// Dictionary corruption would produce a non-string key; skip rather than crashing on key.includes().
if (typeof key !== 'string') continue
if (
valueAsIsKeys.has(key) ||
(key === 'message' && inx in resource && !Number.isInteger(Number(resource[inx]))) ||
Expand Down
66 changes: 48 additions & 18 deletions backend/src/lib/server-side-events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import { constants } from 'node:http2'
import type { Transform } from 'node:stream'
import { clearInterval } from 'node:timers'
import type { Zlib } from 'node:zlib'
import { batchPromiseAll } from './batch-promise-all'
import { getEncodeStream, inflateEvent } from './compression'
import { setCookie } from './cookies'
import { logger } from './logger'
Expand All @@ -17,7 +16,7 @@ import { sizeOf } from '../routes/aggregators/utils'

// If a client hasn't finished receiving a broadcast in PURGE_CLIENT_TIMEOUT
// assume the browser has been refreshed or closed
const PURGE_CLIENT_TIMEOUT = 4 * 60 * 60 * 1000
const PURGE_CLIENT_TIMEOUT = 30 * 60 * 1000

@Ginxo Ginxo Sep 16, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍 aligned to 5.1


const instanceID = randomString(8)

Expand All @@ -35,6 +34,15 @@ export interface ServerSideEvent<DataT = unknown> {
namespace?: string
data?: DataT
}

/** Lightweight resource identity for RBAC filtering without inflating compressed objects. */
export interface EventResourceMeta {
kind: string
apiVersion: string
name?: string
namespace?: string
}

export interface WatchEvent {
type: 'ADDED' | 'DELETED' | 'MODIFIED' | 'EOP'
object: {
Expand All @@ -46,6 +54,24 @@ export interface WatchEvent {
resourceVersion: string
}
}
meta?: EventResourceMeta
}

/** Resolve kind/apiVersion/name/namespace from meta or an already-inflated object. */
export function getEventResourceMeta(event: ServerSideEvent): EventResourceMeta | undefined {
const data = event.data as (WatchEvent & { type?: string }) | undefined
if (!data || typeof data !== 'object') return undefined
if (data.meta?.kind) return data.meta
const object = data.object as WatchEvent['object'] | Buffer | undefined
if (object && !Buffer.isBuffer(object) && typeof object === 'object' && object.kind) {
return {
kind: object.kind,
apiVersion: object.apiVersion,
name: object.metadata?.name,
namespace: object.metadata?.namespace,
}
}
return undefined
}

export interface ServerSideEventClient {
Expand Down Expand Up @@ -125,22 +151,26 @@ export class ServerSideEvents {
}
}

private static async sendEvent(clientID: string, event: ServerSideEvent): Promise<void> {
private static sendEvent(clientID: string, event: ServerSideEvent): Promise<void> {
const client = this.clients[clientID]
if (!client) return
if (client.events && !client.events[event.name]) return
if (client.namespaces && !client.namespaces[event.namespace]) return
event = await inflateEvent(event)
if (!client) return Promise.resolve()
if (client.events && !client.events[event.name]) return Promise.resolve()
if (client.namespaces && !client.namespaces[event.namespace]) return Promise.resolve()
// Filter before inflate so denied events never materialize full resource JSON in memory.
if (this.eventFilter) {
client.eventQueue.push(
this.eventFilter(client.token, event)
.then((shouldSendEvent) => (shouldSendEvent ? event : undefined))
.then((shouldSendEvent) => {
if (!shouldSendEvent) return undefined
return inflateEvent(event)
})
.catch((): undefined => undefined)
)
} else {
client.eventQueue.push(Promise.resolve(event))
client.eventQueue.push(inflateEvent(event))
}
void this.processClient(clientID)
return Promise.resolve()
}

private static async processClient(clientID: string): Promise<void> {
Expand Down Expand Up @@ -243,8 +273,8 @@ export class ServerSideEvents {
res: Http2ServerResponse
): Promise<ServerSideEventClient> {
const [writableStream, compressionStream, encoding] = getEncodeStream(
res as unknown as NodeJS.WritableStream,
req.headers[HTTP2_HEADER_ACCEPT_ENCODING] as string,
res,
req.headers[HTTP2_HEADER_ACCEPT_ENCODING],
process.env.DISABLE_STREAM_COMPRESSION === 'true'
)

Expand Down Expand Up @@ -311,10 +341,10 @@ export class ServerSideEvents {

// SORT EVENTS INTO SMALLER PACKETS
// SO THAT BROWSER PAGE LOADS QUICKER
// uncompress and split events into packets
// Classify using meta / inflated object identity — do not inflate the whole cache up front.
const values = Object.values(this.events)
const compressed = sizeOf(values)
let parts = await batchPromiseAll(values, (event) => inflateEvent(event))
let parts: ServerSideEvent[] = [...values]

// mock a large environment
if (process.env.MOCK_CLUSTERS) {
Expand Down Expand Up @@ -343,9 +373,9 @@ export class ServerSideEvents {
const other: ServerSideEvent<unknown>[] = []
const remainder: ServerSideEvent<unknown>[] = []
parts.forEach((event) => {
const data = event.data as WatchEvent
const meta = getEventResourceMeta(event)
// see frontend/src/components/LoadPluginData.tsx for what pages are fast loaded
switch (data.object.kind) {
switch (meta?.kind) {
case 'ManagedCluster':
case 'HostedCluster':
case 'ClusterDeployment':
Expand Down Expand Up @@ -384,9 +414,9 @@ export class ServerSideEvents {
// sort events alphabetically so that browser list fills from top to bottom
const compareFn =
(propName: 'name' | 'namespace') => (a: ServerSideEvent<unknown>, b: ServerSideEvent<unknown>) => {
const adata = a.data as WatchEvent
const bdata = b.data as WatchEvent
return adata.object.metadata[propName].localeCompare(bdata.object.metadata[propName])
const aVal = getEventResourceMeta(a)?.[propName] ?? ''
const bVal = getEventResourceMeta(b)?.[propName] ?? ''
return aVal.localeCompare(bVal)
}
clusters.sort(compareFn('name'))
infos.sort(compareFn('namespace'))
Expand Down
6 changes: 6 additions & 0 deletions backend/src/resources/watch-options.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,4 +7,10 @@ export interface IWatchOptions {
// poll the resource list instead of watching it
// process the items in its own cache so not to overload event cache
isPolled?: boolean
/**
* True when the Kubernetes resource is cluster-scoped.
* Used by SSE RBAC to decide whether SelfSubjectRulesReview should probe `default`
* (cluster-scoped) or the resource namespace (namespaced).
*/
clusterScoped?: boolean
}
Loading