import * as Lambda from "@distilled.cloud/aws/lambda";
import type * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Schedule from "effect/Schedule";
import { deepEqual } from "../../Diff.ts";
import { toWireSeconds } from "../../Util/Duration.ts";
const DEFAULT_MAXIMUM_RETRY_ATTEMPTS = 2;
const DEFAULT_MAXIMUM_EVENT_AGE_SECONDS = 21_600;
/**
* Asynchronous invocation settings for a Lambda function or alias.
*
* Configured via {@link FunctionProps.eventInvokeConfig} for the unqualified
* function, or {@link AliasProps.eventInvokeConfig} for a specific alias.
*/
export interface EventInvokeConfig {
/**
* Maximum number of times Lambda retries an asynchronous invocation.
* @default 2
*/
maximumRetryAttempts?: number;
/**
* Maximum age that Lambda retains an asynchronous event (e.g. `"6 hours"`).
* Rounded to whole seconds on the wire.
* @default "6 hours"
*/
maximumEventAge?: Duration.Input;
/**
* Destinations for successful or failed asynchronous invocation records.
*/
destinationConfig?: Lambda.DestinationConfig;
}
const retryOnConflict = (
effect: Effect.Effect,
) =>
effect.pipe(
Effect.retry({
while: (e) => e._tag === "ResourceConflictException",
schedule: Schedule.max([Schedule.exponential(500), Schedule.recurs(10)]),
}),
);
const normalizeDestinationConfig = (
config: Lambda.DestinationConfig | undefined,
): Lambda.DestinationConfig | undefined => {
const onSuccess = config?.OnSuccess?.Destination
? { Destination: config.OnSuccess.Destination }
: undefined;
const onFailure = config?.OnFailure?.Destination
? { Destination: config.OnFailure.Destination }
: undefined;
return onSuccess || onFailure
? {
OnSuccess: onSuccess,
OnFailure: onFailure,
}
: undefined;
};
const observeConfig = (functionName: string, qualifier: string | undefined) =>
Lambda.getFunctionEventInvokeConfig({
FunctionName: functionName,
Qualifier: qualifier,
}).pipe(
Effect.catchTag("ResourceNotFoundException", () =>
Effect.succeed(undefined),
),
);
/**
* Converge the async invocation config of a function or alias to the desired
* state: put it when set, delete it when omitted. Diffs against the observed
* cloud config so a no-op deploy skips the write entirely.
*/
export const syncEventInvokeConfig = Effect.fn(function* ({
functionName,
qualifier,
config,
}: {
functionName: string;
/** Alias name (or version) to scope the config to. Omit for `$LATEST`. */
qualifier?: string;
config: EventInvokeConfig | undefined;
}) {
const observed = yield* observeConfig(functionName, qualifier);
if (config === undefined) {
if (observed) {
yield* retryOnConflict(
Lambda.deleteFunctionEventInvokeConfig({
FunctionName: functionName,
Qualifier: qualifier,
}),
).pipe(Effect.catchTag("ResourceNotFoundException", () => Effect.void));
}
return;
}
const desired = {
maximumRetryAttempts:
config.maximumRetryAttempts ?? DEFAULT_MAXIMUM_RETRY_ATTEMPTS,
maximumEventAgeSeconds:
toWireSeconds(config.maximumEventAge) ??
DEFAULT_MAXIMUM_EVENT_AGE_SECONDS,
destinationConfig: normalizeDestinationConfig(config.destinationConfig),
};
if (
observed &&
(observed.MaximumRetryAttempts ?? DEFAULT_MAXIMUM_RETRY_ATTEMPTS) ===
desired.maximumRetryAttempts &&
(observed.MaximumEventAgeInSeconds ?? DEFAULT_MAXIMUM_EVENT_AGE_SECONDS) ===
desired.maximumEventAgeSeconds &&
deepEqual(
normalizeDestinationConfig(observed.DestinationConfig),
desired.destinationConfig,
)
) {
return;
}
yield* Lambda.putFunctionEventInvokeConfig({
FunctionName: functionName,
Qualifier: qualifier,
MaximumRetryAttempts: desired.maximumRetryAttempts,
MaximumEventAgeInSeconds: desired.maximumEventAgeSeconds,
DestinationConfig: desired.destinationConfig,
}).pipe(
Effect.retry({
while: (
e,
): e is
| Lambda.ResourceConflictException
| Lambda.InvalidParameterValueException =>
e._tag === "ResourceConflictException" ||
// Destination validation races IAM policy propagation on the
// execution role — Lambda rejects the put until the role can reach
// the destination.
(e._tag === "InvalidParameterValueException" &&
(e.message?.includes(
"The function execution role does not have permissions to call",
) ??
false)),
schedule: Schedule.max([Schedule.exponential(500), Schedule.recurs(10)]),
}),
);
});