Skip to content
Merged
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
34 changes: 26 additions & 8 deletions packages/@aws-cdk/toolkit-lib/lib/api/logs-monitor/logs-monitor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ import * as util from 'node:util';
import type * as cxapi from '@aws-cdk/cloud-assembly-api';
import chalk from 'chalk';
import type { CloudWatchLogEvent } from '../../payloads/logs-monitor';
import { flatten } from '../../util';
import type { SDK } from '../aws-auth/private';
import { IO } from '../io/private';
import type { IoHelper } from '../io/private';
Expand Down Expand Up @@ -156,31 +155,50 @@ export class CloudWatchLogEventMonitor {
/* c8 ignore stop */

try {
const events = flatten(await this.readNewEvents());
for (const event of events) {
await this.print(event);
for (const result of await this.readNewEvents()) {
if ('error' in result) {
await this.reportError(result.error);
continue;
}
for (const event of result.events) {
await this.print(event);
}
}

// We might have been stop()ped while the network call was in progress.
if (!this.monitorId) {
return;
}
} catch (e: any) {
await this.ioHelper.notify(IO.CDK_TOOLKIT_E5035.msg(`Error occurred while monitoring logs: ${String(e)}`, { error: e }));
await this.reportError(e);
}

this.scheduleNextTick();
}

private async reportError(e: any): Promise<void> {
if (e?.name === 'ThrottlingException') {
// The log group keeps its start time, so the events are read on the next tick.
await this.ioHelper.defaults.debug(`Throttled while monitoring logs, will retry: ${String(e)}`);
return;
}
await this.ioHelper.notify(IO.CDK_TOOLKIT_E5035.msg(`Error occurred while monitoring logs: ${String(e)}`, { error: e }));
}

/**
* Reads all new log events from a set of CloudWatch Log Groups
* in parallel
*
* A failure for one log group does not discard the events of the other log groups.
*/
private async readNewEvents(): Promise<Array<Array<CloudWatchLogEvent>>> {
const promises: Array<Promise<Array<CloudWatchLogEvent>>> = [];
private async readNewEvents(): Promise<Array<{ events: CloudWatchLogEvent[] } | { error: any }>> {
const promises: Array<Promise<{ events: CloudWatchLogEvent[] } | { error: any }>> = [];
for (const settings of this.envsLogGroupsAccessSettings.values()) {
for (const group of Object.keys(settings.logGroupsStartTimes)) {
promises.push(this.readEventsFromLogGroup(settings, group));
promises.push(this.readEventsFromLogGroup(settings, group).then(
(events) => ({ events }),
(error) => ({ error }),
));
}
}
// Limited set of log groups
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,60 @@ test('process truncated events', async () => {
);
});

describe('when reading one log group fails', () => {
const env = { name: 'name', account: '11111111111', region: 'us-east-1' };
const throttle = Object.assign(new Error('Rate exceeded'), { name: 'ThrottlingException' });

function messages(): string[] {
return ioHost.notifySpy.mock.calls.map((call) => stripAnsi(call[0].message));
}

test('a throttled log group does not discard the events of the other log groups, and is read again later', async () => {
// GIVEN
mockCloudWatchClient.on(FilterLogEventsCommand, { logGroupName: 'group-a' })
.resolvesOnce({ events: [event(102, 'message-a', new Date(T102))] })
.resolves({ events: [] });
mockCloudWatchClient.on(FilterLogEventsCommand, { logGroupName: 'group-b' })
.rejectsOnce(throttle)
.resolvesOnce({ events: [event(103, 'message-b', new Date(T102))] })
.resolves({ events: [] });
monitor.addLogGroups(env, sdk, ['group-a', 'group-b']);

// WHEN
await monitor.activate();
await sleep(2500);

// THEN
expect(messages()).toEqual([
expect.stringContaining('[group-a]'),
expect.stringContaining('[group-b]'),
]);
expect(messages()[0]).toContain('message-a');
expect(messages()[1]).toContain('message-b');
});

test('other errors are still reported, without discarding the events of the other log groups', async () => {
// GIVEN
mockCloudWatchClient.on(FilterLogEventsCommand, { logGroupName: 'group-a' })
.resolvesOnce({ events: [event(102, 'message-a', new Date(T102))] })
.resolves({ events: [] });
mockCloudWatchClient.on(FilterLogEventsCommand, { logGroupName: 'group-b' })
.rejectsOnce(new Error('Access denied'))
.resolves({ events: [] });
monitor.addLogGroups(env, sdk, ['group-a', 'group-b']);

// WHEN
await monitor.activate();
await sleep(500);

// THEN
expect(messages()).toEqual([
expect.stringContaining('message-a'),
expect.stringContaining('Error occurred while monitoring logs: Error: Access denied'),
]);
});
});

const T0 = 1597837230504;
const T100 = T0 + 100 * 1000;
const T102 = T0 + 102 * 1000;
Expand Down
Loading