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
2 changes: 1 addition & 1 deletion packages/common-aws/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@paradoxical-io/common-aws",
"version": "2.8.0",
"version": "2.9.0",
"description": "Common aws for paradox services",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
Expand Down
22 changes: 18 additions & 4 deletions packages/common-aws/src/cloudwatch/cloudwatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,27 @@ import {
FilterLogEventsCommand,
LogGroup,
} from '@aws-sdk/client-cloudwatch-logs';
import { log } from '@paradoxical-io/common-server';
import { Brand, EpochMS } from '@paradoxical-io/types';

import { Logger, Monitoring, noOpMonitoring } from '../monitoring';

export type CloudwatchExportTaskId = Brand<'TaskId', string>;

export class CloudwatchManager {
constructor(private cloudwatch: CloudWatchLogsClient = new CloudWatchLogsClient()) {}
private readonly logger: Logger;

constructor({
cloudwatch = new CloudWatchLogsClient(),
monitoring = noOpMonitoring(),
}: {
cloudwatch?: CloudWatchLogsClient;
monitoring?: Monitoring;
} = {}) {
this.cloudwatch = cloudwatch;
this.logger = monitoring.logger;
}

private readonly cloudwatch: CloudWatchLogsClient;

/**
* Create a new task to export a log group to an S3 bucket.
Expand All @@ -36,7 +50,7 @@ export class CloudwatchManager {
from: EpochMS;
to: EpochMS;
}): Promise<CloudwatchExportTaskId | undefined> {
log.info(`Creating export task for ${logGroupName} from ${from} to ${to}`);
this.logger.info(`Creating export task for ${logGroupName} from ${from} to ${to}`);

/**
* validate that we have at least 1 log event to export within the time range.
Expand Down Expand Up @@ -66,7 +80,7 @@ export class CloudwatchManager {
return taskId as CloudwatchExportTaskId;
}

log.info(`No log events found for ${logGroupName} from ${from} to ${to}`);
this.logger.info(`No log events found for ${logGroupName} from ${from} to ${to}`);

return undefined;
}
Expand Down
11 changes: 7 additions & 4 deletions packages/common-aws/src/credentials.ts
Original file line number Diff line number Diff line change
@@ -1,19 +1,22 @@
import { gitRootSync, log } from '@paradoxical-io/common-server';
import { gitRootSync } from '@paradoxical-io/common-server';
import { Env } from '@paradoxical-io/types';
import * as path from 'path';

import { Logger, noOpMonitoring } from './monitoring';

export interface Profile {
profile: string;
}

/**
* Sets the current environment to use the AWS credentials from the shared credentials file
* @param requestedEnv
* @param logger
*/
export function useAWS(requestedEnv: Env | Profile): void {
export function useAWS(requestedEnv: Env | Profile, logger: Logger = noOpMonitoring().logger): void {
let env = requestedEnv;
if (env === 'local') {
log.debug("Running 'local' environment. Using dev env for all aws resources");
logger.debug("Running 'local' environment. Using dev env for all aws resources");
env = 'dev';
}

Expand All @@ -24,7 +27,7 @@ export function useAWS(requestedEnv: Env | Profile): void {
process.env.AWS_CONFIG_FILE = configPath;
process.env.AWS_SDK_LOAD_CONFIG = 'true';

log.debug(`Configured AWS to use profile ${profile}`);
logger.debug(`Configured AWS to use profile ${profile}`);
}

function resolveConfig(env: Env | Profile) {
Expand Down
15 changes: 12 additions & 3 deletions packages/common-aws/src/dynamo/do-once/doOnce.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,20 @@
import { defaultTimeProvider } from '@paradoxical-io/common';
import { logMethod } from '@paradoxical-io/common-server';
import { defaultTimeProvider, logMethod } from '@paradoxical-io/common';
import { CompoundKey, DoOnceActionKey, DoOnceResponse, EpochMS, SortKey } from '@paradoxical-io/types';

import { Logger, Monitoring, noOpMonitoring } from '../../monitoring';
import { PartitionedKeyValueTable } from '../keys';

export class DoOnceManager<Key extends string = string> {
constructor(private readonly kv: PartitionedKeyValueTable, private readonly time = defaultTimeProvider()) {}
// Accessed by @logMethod() decorator via reflection
readonly logger: Logger;

constructor(
private readonly kv: PartitionedKeyValueTable,
private readonly time = defaultTimeProvider(),
monitoring: Monitoring = noOpMonitoring()
) {
this.logger = monitoring.logger;
}

/**
* The key format where each action is stored separately in the KV.
Expand Down
15 changes: 11 additions & 4 deletions packages/common-aws/src/dynamo/dynamoLock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,10 @@ import {
UpdateItemCommand,
} from '@aws-sdk/client-dynamodb';
import { defaultTimeProvider, propertiesOf, TimeProvider, toEpochSeconds } from '@paradoxical-io/common';
import { Lock, LockApi, log } from '@paradoxical-io/common-server';
import { Lock, LockApi } from '@paradoxical-io/common-server';
import { EpochSeconds } from '@paradoxical-io/types';

import { Logger, Monitoring, noOpMonitoring } from '../monitoring';
import { DynamoDao } from './mapper';
import { DynamoTableName, dynamoTableName } from './util';

Expand All @@ -19,20 +20,26 @@ export class DynamoLock implements LockApi {

private readonly timeProvider: TimeProvider;

private readonly logger: Logger;

constructor({
dynamo = new DynamoDBClient(),
tableName = dynamoTableName('locks'),
timeProvider = defaultTimeProvider(),
monitoring = noOpMonitoring(),
}: {
dynamo?: DynamoDBClient;
tableName?: DynamoTableName;
timeProvider?: TimeProvider;
monitoring?: Monitoring;
} = {}) {
this.dynamo = dynamo;

this.tableName = tableName;

this.timeProvider = timeProvider;

this.logger = monitoring.logger;
}

async tryAcquire(key: string, timeoutSeconds: number): Promise<Lock | undefined> {
Expand Down Expand Up @@ -79,7 +86,7 @@ export class DynamoLock implements LockApi {
expired = `because previous lock expired at ${old.Attributes['expiresAt']!.N}`;
}

log.info(
this.logger.info(
`Acquired lock id: ${key}, expires at ${payload.expiresAt}, evaluated with now as ${now.toString()} ${expired}`
);

Expand All @@ -97,12 +104,12 @@ export class DynamoLock implements LockApi {

await this.dynamo.send(deleteItemCommand);

log.info(`Released lock id: ${key}`);
this.logger.info(`Released lock id: ${key}`);
},
};
} catch (e) {
if (e instanceof ConditionalCheckFailedException) {
log.info(`Cannot acquire lock id: ${key}, it is still held`);
this.logger.info(`Cannot acquire lock id: ${key}, it is still held`);

return undefined;
}
Expand Down
16 changes: 11 additions & 5 deletions packages/common-aws/src/dynamo/keys/keyTable.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,9 @@ import {
ScanCommandInput,
} from '@aws-sdk/client-dynamodb';
import { Arrays, asBrandSafe, propertyOf, sleep } from '@paradoxical-io/common';
import { log } from '@paradoxical-io/common-server';
import { Brand, Milliseconds, notNullOrUndefined } from '@paradoxical-io/types';

import { Logger, Monitoring, noOpMonitoring } from '../../monitoring';
import { DynamoDao } from '../mapper';
import { assertTableNameValid, DynamoTableName, dynamoTableName } from '../util';

Expand Down Expand Up @@ -54,14 +54,18 @@ export class KeyValueTable {

private readonly tableName: string;

private readonly logger: Logger;

constructor({
namespace = 'global',
dynamo = new DynamoDBClient(),
tableName = dynamoTableName('keys'),
monitoring = noOpMonitoring(),
}: {
namespace?: string;
dynamo?: DynamoDBClient;
tableName?: DynamoTableName;
monitoring?: Monitoring;
} = {}) {
assertTableNameValid(tableName);

Expand All @@ -70,6 +74,8 @@ export class KeyValueTable {
this.tableName = tableName;

this.dynamo = dynamo;

this.logger = monitoring.logger;
}

async get<T>(id: string): Promise<T | undefined> {
Expand Down Expand Up @@ -168,7 +174,7 @@ export class KeyValueTable {
}
}
} catch (e) {
log.error(`Failed to scan ${pageItem}`, e);
this.logger.error(`Failed to scan ${pageItem}`, e);

throw e;
}
Expand Down Expand Up @@ -229,7 +235,7 @@ export class KeyValueTable {
}
}
} catch (e) {
log.error(`Failed to get batch ${id.join(',')}`, e);
this.logger.error(`Failed to get batch ${id.join(',')}`, e);

throw e;
}
Expand All @@ -254,7 +260,7 @@ export class KeyValueTable {
try {
await this.dynamo.send(command);
} catch (e) {
log.error(`Failed to set item ${data.key}`, e);
this.logger.error(`Failed to set item ${data.key}`, e);

throw e;
}
Expand All @@ -273,7 +279,7 @@ export class KeyValueTable {
try {
await this.dynamo.send(command);
} catch (e) {
log.error(`Failed to delete item ${id}`, e);
this.logger.error(`Failed to delete item ${id}`, e);

throw e;
}
Expand Down
12 changes: 9 additions & 3 deletions packages/common-aws/src/dynamo/keys/keyValueCounter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,9 @@ import {
UpdateItemCommand,
} from '@aws-sdk/client-dynamodb';
import { Arrays, propertyOf } from '@paradoxical-io/common';
import { log } from '@paradoxical-io/common-server';
import { Brand, notNullOrUndefined } from '@paradoxical-io/types';

import { Logger, Monitoring, noOpMonitoring } from '../../monitoring';
import { DynamoDao } from '../mapper';
import { assertTableNameValid, DynamoTableName, dynamoTableName } from '../util';

Expand Down Expand Up @@ -81,14 +81,18 @@ export class KeyValueCounter {

private readonly tableName: string;

private readonly logger: Logger;

constructor({
namespace,
dynamo = new DynamoDBClient(),
tableName = dynamoTableName('keys'),
monitoring = noOpMonitoring(),
}: {
namespace: string;
dynamo?: DynamoDBClient;
tableName?: DynamoTableName;
monitoring?: Monitoring;
}) {
assertTableNameValid(tableName);

Expand All @@ -97,6 +101,8 @@ export class KeyValueCounter {
this.tableName = tableName;

this.dynamo = dynamo;

this.logger = monitoring.logger;
}

async get<K extends string = string>(ids: K[]): Promise<Array<KeyCount<K>>> {
Expand Down Expand Up @@ -217,7 +223,7 @@ export class KeyValueCounter {

return Number(result.Attributes![propertyOf<KeyValueCountTableDao>('count')]!.N);
} catch (e) {
log.error(`Failed to increment ${data.key}`, e);
this.logger.error(`Failed to increment ${data.key}`, e);

throw e;
}
Expand All @@ -236,7 +242,7 @@ export class KeyValueCounter {
try {
await this.dynamo.send(command);
} catch (e) {
log.error(`Failed to delete ${id}`, e);
this.logger.error(`Failed to delete ${id}`, e);

throw e;
}
Expand Down
Loading
Loading