Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions backend/src/lib/multi-cluster-hub.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
/* Copyright Contributors to the Open Cluster Management project */

import { getServiceAccountToken } from './serviceAccountToken'
import { jsonRequest } from './json-request'
import { logger } from './logger'
import { getServiceAccountToken } from './serviceAccountToken'

// Type returned by /apis/authentication.k8s.io/v1/tokenreviews

Expand All @@ -25,7 +25,7 @@ interface MultiClusterHubList {
items: MultiClusterHub[]
}

let multiclusterhub: Promise<MultiClusterHub | undefined>
let multiclusterhub: Promise<MultiClusterHub | undefined> | undefined

/** Clear MultiClusterHub cache. Used for test isolation. */
export function resetMultiClusterHubCache(): void {
Expand All @@ -40,10 +40,10 @@ export async function getMultiClusterHub(noCache?: boolean): Promise<MultiCluste
serviceAccountToken
)
.then((response) => {
return response.items && response.items[0] ? response.items[0] : undefined
return response.items?.[0] ?? undefined
})
.catch((err: Error): undefined => {
logger.error({ msg: 'Error getting MultiClusterHub', error: err.message })
logger.debug({ msg: 'MultiClusterHub not found', error: err.message })
return undefined
})
}
Expand All @@ -52,5 +52,5 @@ export async function getMultiClusterHub(noCache?: boolean): Promise<MultiCluste

export async function getMultiClusterHubComponents(noCache?: boolean): Promise<MultiClusterHubComponent[] | undefined> {
const multiClusterHub = await getMultiClusterHub(noCache)
return multiClusterHub.spec?.overrides?.components
return multiClusterHub?.spec?.overrides?.components
}
155 changes: 96 additions & 59 deletions backend/src/lib/search.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
/* Copyright Contributors to the Open Cluster Management project */
import type { IncomingMessage } from 'node:http'
import type { OutgoingHttpHeaders } from 'node:http2'
import type { RequestOptions } from 'node:https'
import { request } from 'node:https'
import { pipeline } from 'node:stream/promises'
import { Writable } from 'node:stream'
import { URL } from 'node:url'
import { getMultiClusterHub } from '../lib/multi-cluster-hub'
import { getNamespace, getServiceAccountToken } from '../lib/serviceAccountToken'
Expand All @@ -26,6 +29,21 @@ export type ISearchResult = {
message?: string
}

function collectResponseBody(res: IncomingMessage): Promise<string> {
return new Promise((resolve, reject) => {
let body = ''
const collector = new Writable({
write(chunk: Buffer, _encoding, callback) {
body += chunk.toString()
callback()
},
})
pipeline(res, collector)
.then(() => resolve(body))
.catch(reject)
})
}

export async function getServiceAccountSearchRequestOptions() {
const serviceAccountToken = getServiceAccountToken()
const headers: OutgoingHttpHeaders = {
Expand All @@ -38,9 +56,10 @@ export async function getServiceAccountSearchRequestOptions() {
}

export async function getSearchRequestOptions(headers: OutgoingHttpHeaders): Promise<RequestOptions> {
const mch = await getMultiClusterHub()
const multiClusterHub = await getMultiClusterHub()
const namespace = getNamespace()
const machineNs = process.env.NODE_ENV === 'test' ? 'undefined' : `${mch?.metadata?.namespace || namespace}`
const machineNs =
process.env.NODE_ENV === 'test' ? 'undefined' : `${multiClusterHub?.metadata?.namespace || namespace}`
const searchService = `https://search-search-api.${machineNs}.svc.cluster.local:4010`
const searchUrl = process.env.SEARCH_API_URL || searchService
const endpoint = process.env.globalSearchFeatureFlag === 'enabled' ? '/federated' : '/searchapi/graphql'
Expand All @@ -62,41 +81,51 @@ export async function getSearchResults(query: IQuery) {
const options = await getServiceAccountSearchRequestOptions()
const requestTimeout = 2 * 60 * 1000
return new Promise<ISearchResult>((resolve, reject) => {
let body = ''
const id = setTimeout(() => {
logger.error(`getSearchResults request timeout`)
reject(new Error('request timeout'))
}, requestTimeout)
const req = request(options, (res) => {
res.on('data', (data) => {
body += data
})
res.on('end', () => {
try {
const result = JSON.parse(body) as ISearchResult
const message = typeof result === 'string' ? result : result.message
if (message) {
logger.error(`getSearchResults return error ${message}`)
reject(new Error(result.message))
let settled = false
const timeout = { requestTimeoutId: undefined as NodeJS.Timeout | undefined }
const finish = (fn: () => void) => {
if (settled) return
settled = true
clearTimeout(timeout.requestTimeoutId)
fn()
}
const clientRequest = request(options, (res) => {
void collectResponseBody(res)
.then((body) => {
try {
const result = JSON.parse(body) as ISearchResult
const message = typeof result === 'string' ? result : result.message
if (message) {
logger.error(`getSearchResults return error ${message}`)
finish(() => reject(new Error(result.message)))
return
}
finish(() => resolve(result))
} catch (e) {
// search might be overwhelmed
// pause before next request
logger.error(`getSearchResults parse error ${e} ${body}`)
clearTimeout(timeout.requestTimeoutId)
setTimeout(() => {
finish(() => reject(new Error(body)))
}, requestTimeout)
}
resolve(result)
} catch (e) {
// search might be overwhelmed
// pause before next request
logger.error(`getSearchResults parse error ${e} ${body}`)
setTimeout(() => {
reject(new Error(body))
}, requestTimeout)
}
clearTimeout(id)
})
})
.catch((e: Error) => {
finish(() => reject(e))
})
})
req.on('error', (e) => {
timeout.requestTimeoutId = setTimeout(() => {
logger.error(`getSearchResults request timeout`)
clientRequest.destroy()
finish(() => reject(new Error('request timeout')))
}, requestTimeout)
clientRequest.on('error', (e) => {
logger.error(`getSearchResults request error ${e.message}`)
reject(e)
finish(() => reject(e))
})
req.write(JSON.stringify(query))
req.end()
clientRequest.write(JSON.stringify(query))
clientRequest.end()
})
}

Expand Down Expand Up @@ -125,37 +154,45 @@ const ping = {
export async function pingSearchAPI() {
const options = await getServiceAccountSearchRequestOptions()
return new Promise<boolean>((resolve, reject) => {
let body = ''
const id = setTimeout(
let settled = false
const timeout = { requestTimeoutId: undefined as NodeJS.Timeout | undefined }
const finish = (fn: () => void) => {
if (settled) return
settled = true
clearTimeout(timeout.requestTimeoutId)
fn()
}
const clientRequest = request(options, (res) => {
void collectResponseBody(res)
.then((body) => {
try {
const result = JSON.parse(body) as { data: unknown }
if (result.data) {
finish(() => resolve(true))
} else {
finish(() => reject(new Error('no data')))
}
} catch (e) {
logger.error(`pingSearchAPI parse error ${e} ${body}`)
finish(() => reject(new Error(String(e).valueOf())))
}
})
.catch((e: Error) => {
finish(() => reject(e))
})
})
timeout.requestTimeoutId = setTimeout(
() => {
logger.error(`ping searchAPI timeout`)
reject(new Error('request timeout'))
clientRequest.destroy()
finish(() => reject(new Error('request timeout')))
},
4 * 60 * 1000
)
const req = request(options, (res) => {
res.on('data', (data) => {
body += data
})
res.on('end', () => {
try {
const result = JSON.parse(body) as { data: unknown }
if (result.data) {
resolve(true)
} else {
reject(new Error('no data'))
}
} catch (e) {
logger.error(`pingSearchAPI parse error ${e} ${body}`)
reject(new Error(String(e).valueOf()))
}
clearTimeout(id)
})
})
req.on('error', (e) => {
reject(e)
clientRequest.on('error', (e) => {
finish(() => reject(e))
})
req.write(JSON.stringify(ping))
req.end()
clientRequest.write(JSON.stringify(ping))
clientRequest.end()
})
}
40 changes: 37 additions & 3 deletions backend/src/routes/aggregators/applications.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { addOCPQueryInputs, addSystemQueryInputs, cacheOCPApplications } from '.
import { ApplicationSetKind, type IApplicationSet, type IResource, type SearchResult } from '../../resources/resource'
import type { FilterSelections, ISortBy } from '../../lib/pagination'
import { logger } from '../../lib/logger'
import { getMultiClusterHub } from '../../lib/multi-cluster-hub'
import {
discoverSystemAppNamespacePrefixes,
getApplicationsHelper,
Expand Down Expand Up @@ -196,12 +197,30 @@ export const promiseTimeout = <T>(promise: Promise<T>, delay: number) => {
// //////////////////////////////////////////////////////////////////////////////////
export async function startAggregatingApplications() {
await discoverSystemAppNamespacePrefixes()
void searchLoop()
await searchLoop()
}

let stopping = false
let cancelPendingWait: (() => void) | undefined

function waitWhileRunning(ms: number): Promise<void> {
if (stopping) return Promise.resolve()
return new Promise((resolve) => {
const timeoutId = setTimeout(() => {
cancelPendingWait = undefined
resolve()
}, ms)
cancelPendingWait = () => {
clearTimeout(timeoutId)
cancelPendingWait = undefined
resolve()
}
})
}

export function stopAggregatingApplications(): void {
stopping = true
cancelPendingWait?.()
}

/** Reset aggregation stopping flag. Used for test isolation. */
Expand Down Expand Up @@ -354,7 +373,22 @@ export async function addUIData(items: ITransformedResource[]) {
export async function searchLoop() {
let pass = 1
let searchAPIMissing = false
let multiClusterHubMissing = false
while (!stopping) {
const multiClusterHub = await getMultiClusterHub(true)
if (!multiClusterHub) {
if (!multiClusterHubMissing) {
logger.info('MultiClusterHub not found; waiting before search aggregation')
multiClusterHubMissing = true
}
await waitWhileRunning(5 * 60 * 1000)
continue
}
if (multiClusterHubMissing) {
logger.info('MultiClusterHub found')
multiClusterHubMissing = false
}

// make sure there's an active search api
// otherwise there's no point
let exists
Expand All @@ -371,7 +405,7 @@ export async function searchLoop() {
logger.error('search API missing')
searchAPIMissing = true
}
await new Promise((r) => setTimeout(r, 5 * 60 * 1000))
await waitWhileRunning(5 * 60 * 1000)
}
} while (!exists)
/* istanbul ignore if */
Expand All @@ -394,7 +428,7 @@ export async function searchLoop() {
// process every APP_SEARCH_INTERVAL seconds
/* istanbul ignore if */
if (process.env.NODE_ENV !== 'test') {
await new Promise((r) => setTimeout(r, pass <= 3 ? 15000 : Number(process.env.APP_SEARCH_INTERVAL) || 60000))
await waitWhileRunning(pass <= 3 ? 15000 : Number(process.env.APP_SEARCH_INTERVAL) || 60000)
} else {
stopping = true
}
Expand Down
Loading