Improve logging recovery, resource scaling, and app deploy logging.

Auto-reconnect Elasticsearch port-forward after cluster or API restarts, poll log status in the UI, and apply storage changes through billing upgrade for all workloads. Add Redis/RabbitMQ PVC resize, Helm ES credentials for Fluent Bit, and fix deploy progress overlay behavior.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
keyhan
2026-05-25 21:41:00 +03:30
parent bf9e827f85
commit dc9830383b
18 changed files with 1068 additions and 414 deletions
+334 -16
View File
@@ -1,7 +1,16 @@
import { Injectable, Logger, ServiceUnavailableException, Inject, forwardRef } from '@nestjs/common';
import {
Injectable,
Logger,
ServiceUnavailableException,
Inject,
forwardRef,
OnModuleInit,
OnModuleDestroy,
} from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import * as k8s from '@kubernetes/client-node';
import * as crypto from 'crypto';
import { ChildProcess, spawn } from 'child_process';
import { ClustersService } from '../clusters/clusters.service';
import { HelmService, LOGGING_HELM_NAMESPACE, LOGGING_HELM_RELEASE } from './helm.service';
@@ -56,12 +65,18 @@ export interface LogStatsResult {
* that all user apps can send logs to via Fluent Bit sidecars.
*/
@Injectable()
export class ElasticsearchService {
export class ElasticsearchService implements OnModuleInit, OnModuleDestroy {
private readonly logger = new Logger(ElasticsearchService.name);
private readonly ES_NAMESPACE = 'logging';
private readonly ES_NAME = 'elasticsearch';
private readonly KIBANA_NAME = 'kibana';
private portForwardChild: ChildProcess | null = null;
private portForwardStartedByUs = false;
private ensureInFlight: Promise<void> | null = null;
private reconnectTimer: ReturnType<typeof setTimeout> | null = null;
private healthCheckTimer: ReturnType<typeof setInterval> | null = null;
private reconnectAttempt = 0;
// Default credentials - should be overridden via env in production
private readonly ELASTIC_PASSWORD: string;
private readonly FLUENTBIT_PASSWORD: string;
@@ -78,6 +93,254 @@ export class ElasticsearchService {
this.KIBANA_SYSTEM_PASSWORD = this.configService.get('elasticsearch.kibanaPassword') || 'Kibana2024!System';
}
async onModuleInit(): Promise<void> {
await this.ensureLocalElasticsearchAccess({ waitForCluster: true });
if (this.shouldAutoPortForward()) {
this.healthCheckTimer = setInterval(() => {
void this.periodicElasticsearchHealthCheck();
}, 30_000);
}
}
onModuleDestroy(): void {
if (this.healthCheckTimer) {
clearInterval(this.healthCheckTimer);
this.healthCheckTimer = null;
}
if (this.reconnectTimer) {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = null;
}
this.stopDevPortForward();
}
private isLoopbackHost(host: string): boolean {
return host === '127.0.0.1' || host === 'localhost' || host === '::1';
}
private shouldAutoPortForward(): boolean {
if (this.configService.get<string>('elasticsearch.autoPortForward') === 'false') {
return false;
}
if (process.env.ELASTICSEARCH_AUTO_PORT_FORWARD === 'false') {
return false;
}
const nodeEnv = process.env.NODE_ENV || 'development';
if (nodeEnv === 'production') {
return false;
}
const host = this.configService.get<string>('elasticsearch.host') || '';
return this.isLoopbackHost(host);
}
private stopDevPortForward(): void {
if (!this.portForwardChild) {
return;
}
const startedByUs = this.portForwardStartedByUs;
const child = this.portForwardChild;
this.portForwardChild = null;
this.portForwardStartedByUs = false;
child.kill('SIGTERM');
if (startedByUs) {
this.logger.log('Stopped Elasticsearch kubectl port-forward');
}
}
private schedulePortForwardReconnect(reason: string): void {
if (!this.shouldAutoPortForward()) {
return;
}
if (this.reconnectTimer) {
return;
}
const delay = Math.min(60_000, 2_000 * Math.pow(2, this.reconnectAttempt));
this.reconnectAttempt += 1;
this.logger.warn(
`Elasticsearch port-forward lost (${reason}). Reconnecting in ${Math.round(delay / 1000)}s…`,
);
this.reconnectTimer = setTimeout(() => {
this.reconnectTimer = null;
void this.ensureLocalElasticsearchAccess().then((ok) => {
if (ok) {
this.reconnectAttempt = 0;
}
});
}, delay);
}
private async periodicElasticsearchHealthCheck(): Promise<void> {
if (!this.shouldAutoPortForward()) {
return;
}
const deployed = await this.isDeployed();
if (!deployed) {
return;
}
if (await this.probeElasticsearch()) {
this.reconnectAttempt = 0;
return;
}
this.logger.debug('Elasticsearch health check failed; restoring tunnel…');
await this.ensureLocalElasticsearchAccess();
}
private async waitForLoggingStack(maxWaitMs = 120_000): Promise<boolean> {
const started = Date.now();
while (Date.now() - started < maxWaitMs) {
if (await this.isDeployed()) {
return true;
}
await new Promise((r) => setTimeout(r, 5_000));
}
return false;
}
private async probeElasticsearch(timeoutMs = 3000): Promise<boolean> {
try {
const conn = this.getConnectionInfo();
const auth = Buffer.from(`${conn.username}:${conn.password}`).toString('base64');
const response = await fetch(`http://${conn.host}:${conn.port}/_cluster/health`, {
headers: { Authorization: `Basic ${auth}` },
signal: AbortSignal.timeout(timeoutMs),
});
return response.ok;
} catch {
return false;
}
}
private async waitForElasticsearch(maxWaitMs = 15_000): Promise<boolean> {
const started = Date.now();
while (Date.now() - started < maxWaitMs) {
if (await this.probeElasticsearch(2000)) {
return true;
}
await new Promise((r) => setTimeout(r, 400));
}
return false;
}
private startDevPortForward(localPort: number): void {
if (this.portForwardChild) {
return;
}
const args = [
'port-forward',
'-n',
this.ES_NAMESPACE,
`svc/${this.ES_NAME}`,
`${localPort}:9200`,
];
this.logger.log(`Starting kubectl ${args.join(' ')} (local log search)`);
const child = spawn('kubectl', args, { stdio: ['ignore', 'pipe', 'pipe'] });
this.portForwardChild = child;
this.portForwardStartedByUs = true;
child.on('exit', (code, signal) => {
const wasOurs = this.portForwardChild === child;
if (wasOurs) {
this.portForwardChild = null;
this.portForwardStartedByUs = false;
}
if (wasOurs) {
const reason =
code !== 0 && code !== null
? `exit code ${code}`
: signal
? `signal ${signal}`
: 'connection closed';
this.schedulePortForwardReconnect(reason);
}
});
child.stderr?.on('data', (chunk: Buffer) => {
const line = chunk.toString().trim();
if (line && !line.includes('Handling connection')) {
this.logger.debug(`kubectl port-forward: ${line}`);
}
});
}
/**
* When the API runs on the host with ELASTICSEARCH_HOST=127.0.0.1, open a tunnel to the cluster.
* Safe to call repeatedly (e.g. after cluster/API restart or port-forward drop).
*/
private async ensureLocalElasticsearchAccess(options?: {
waitForCluster?: boolean;
}): Promise<boolean> {
if (this.ensureInFlight) {
await this.ensureInFlight;
return this.probeElasticsearch();
}
this.ensureInFlight = this.ensureLocalElasticsearchAccessImpl(options);
try {
await this.ensureInFlight;
return this.probeElasticsearch();
} finally {
this.ensureInFlight = null;
}
}
private async ensureLocalElasticsearchAccessImpl(options?: {
waitForCluster?: boolean;
}): Promise<void> {
if (!this.shouldAutoPortForward()) {
return;
}
if (await this.probeElasticsearch()) {
this.reconnectAttempt = 0;
return;
}
let deployed = await this.isDeployed();
if (!deployed && options?.waitForCluster) {
this.logger.log('Waiting for logging stack after cluster reconnect…');
deployed = await this.waitForLoggingStack();
}
if (!deployed) {
return;
}
const port = this.configService.get<number>('elasticsearch.port') || 9200;
// Stale tunnel after sleep/reboot: port may be bound but ES unreachable
if (this.portForwardChild) {
this.stopDevPortForward();
await new Promise((r) => setTimeout(r, 300));
}
this.startDevPortForward(port);
const ready = await this.waitForElasticsearch(90_000);
if (ready) {
this.reconnectAttempt = 0;
this.logger.log(`Elasticsearch reachable at 127.0.0.1:${port}`);
} else {
this.stopDevPortForward();
this.logger.warn(
`Could not reach Elasticsearch on 127.0.0.1:${port}. Will retry. Manual: kubectl port-forward -n ${this.ES_NAMESPACE} svc/${this.ES_NAME} ${port}:9200`,
);
this.schedulePortForwardReconnect('probe timeout');
}
}
private localElasticsearchHint(): string {
const conn = this.getConnectionInfo();
if (this.isLoopbackHost(conn.host)) {
return (
`Ensure port ${conn.port} is forwarded to the cluster (the API auto-starts kubectl port-forward in development). ` +
`Manual: kubectl port-forward -n ${this.ES_NAMESPACE} svc/${this.ES_NAME} ${conn.port}:9200`
);
}
if (conn.host.includes('svc.cluster.local') || conn.host.includes('.cluster.')) {
return (
'Run the API inside the cluster, or set ELASTICSEARCH_HOST=127.0.0.1 and keep port-forward running: ' +
`kubectl port-forward -n ${this.ES_NAMESPACE} svc/${this.ES_NAME} ${conn.port}:9200`
);
}
return `Ensure Elasticsearch is listening on ${conn.host}:${conn.port}.`;
}
private async getK8sClients(clusterId?: string) {
const cluster = clusterId
? await this.clustersService.findOne(clusterId)
@@ -178,6 +441,10 @@ export class ElasticsearchService {
elasticPassword: this.ELASTIC_PASSWORD,
fluentbitPassword: this.FLUENTBIT_PASSWORD,
kibanaSystemPassword: this.KIBANA_SYSTEM_PASSWORD,
images: {
elasticsearch: this.configService.get<string>('elasticsearch.images.elasticsearch'),
kibana: this.configService.get<string>('elasticsearch.images.kibana'),
},
});
this.logger.log(
@@ -214,8 +481,10 @@ export class ElasticsearchService {
*/
getConnectionInfo(): { host: string; port: number; username: string; password: string } {
return {
host: `${this.ES_NAME}.${this.ES_NAMESPACE}.svc.cluster.local`,
port: 9200,
host:
this.configService.get<string>('elasticsearch.host') ||
`${this.ES_NAME}.${this.ES_NAMESPACE}.svc.cluster.local`,
port: this.configService.get<number>('elasticsearch.port') || 9200,
username: 'elastic',
password: this.ELASTIC_PASSWORD,
};
@@ -363,6 +632,18 @@ export class ElasticsearchService {
return `logs-user-${userId.split('-')[0]}-*`;
}
private elasticsearchFetch(url: string, auth: string, body: unknown): Promise<Response> {
return fetch(url, {
method: body === undefined ? 'GET' : 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Basic ${auth}`,
},
body: body === undefined ? undefined : JSON.stringify(body),
signal: AbortSignal.timeout(15_000),
});
}
private async esRequest(path: string, body: unknown, clusterId?: string): Promise<any> {
const deployed = await this.isDeployed(clusterId);
if (!deployed) {
@@ -375,14 +656,23 @@ export class ElasticsearchService {
const url = `http://${conn.host}:${conn.port}${path}`;
const auth = Buffer.from(`${conn.username}:${conn.password}`).toString('base64');
const response = await fetch(url, {
method: body === undefined ? 'GET' : 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Basic ${auth}`,
},
body: body === undefined ? undefined : JSON.stringify(body),
});
let response: Response | undefined;
try {
response = await this.elasticsearchFetch(url, auth, body);
} catch (err: any) {
if (this.shouldAutoPortForward()) {
await this.ensureLocalElasticsearchAccess();
try {
response = await this.elasticsearchFetch(url, auth, body);
} catch {
// retry failed
}
}
if (!response) {
this.logger.warn(`Elasticsearch unreachable at ${conn.host}:${conn.port}: ${err?.message || err}`);
throw new ServiceUnavailableException(`Cannot reach Elasticsearch. ${this.localElasticsearchHint()}`);
}
}
if (!response.ok) {
const text = await response.text();
@@ -518,8 +808,36 @@ export class ElasticsearchService {
return (result.hits?.hits || []).map((h: any) => this.normalizeHit(h));
}
async getLoggingStatus(clusterId?: string): Promise<{ available: boolean; deployed: boolean }> {
const deployed = await this.isDeployed(clusterId);
return { available: deployed, deployed };
async getLoggingStatus(
clusterId?: string,
): Promise<{ available: boolean; deployed: boolean; recovering?: boolean; message?: string }> {
let deployed = await this.isDeployed(clusterId);
if (!deployed && this.shouldAutoPortForward()) {
deployed = await this.waitForLoggingStack(8_000);
}
if (!deployed) {
return {
available: false,
deployed: false,
message: 'Central logging is not deployed. Ask an administrator to deploy Elasticsearch.',
};
}
if (this.shouldAutoPortForward() && !(await this.probeElasticsearch())) {
void this.ensureLocalElasticsearchAccess();
}
if (await this.probeElasticsearch()) {
return { available: true, deployed: true };
}
return {
available: false,
deployed: true,
recovering: this.shouldAutoPortForward(),
message: this.shouldAutoPortForward()
? 'Reconnecting to Elasticsearch after cluster or API restart. This usually takes under a minute.'
: `Elasticsearch is running in the cluster, but this backend cannot reach it. ${this.localElasticsearchHint()}`,
};
}
}