/*! * Copyright (c) Microsoft Corporation and contributors. All rights reserved. * Licensed under the MIT License. */ import { TypedEventEmitter } from "@fluidframework/common-utils"; import { ISequencedDocumentAugmentedMessage } from "@fluidframework/protocol-definitions"; import { IContext, IControlMessage, IDeliState, IPartitionLambda, IProducer, IServiceConfiguration, NackMessagesType, IQueuedMessage, INackMessagesControlMessageContents, LambdaCloseType, ISequencedSignalClient, IClientManager, ICheckpointService } from "@fluidframework/server-services-core"; import { Lumber, LumberEventName } from "@fluidframework/server-services-telemetry"; import { IEvent } from "../events"; import { IDeliCheckpointManager } from "./checkpointManager"; /** * @internal */ export declare enum OpEventType { /** * There have been no sequenced ops for X milliseconds since the last message. */ Idle = 0, /** * More than X amount of ops have been ticketed since the emit. */ MaxOps = 1, /** * There was no previous emit for the last X milliseconds. */ MaxTime = 2, /** * Indicates the durable sequence number was updated. */ UpdatedDurableSequenceNumber = 3 } /** * @internal */ export interface IDeliLambdaEvents extends IEvent { /** * Emitted when certain op event heuristics are triggered. */ (event: "opEvent", listener: (type: OpEventType, sequenceNumber: number, sequencedMessagesSinceLastOpEvent: number) => void): any; /** * Emitted when the lambda is updating the durable sequence number. * This usually occurs via a control message after a summary was created. */ (event: "updatedDurableSequenceNumber", listener: (durableSequenceNumber: number) => void): any; /** * Emitted when the lambda is updating a nack message */ (event: "updatedNackMessages", listener: (type: NackMessagesType, contents: INackMessagesControlMessageContents | undefined) => void): any; /** * Emitted when the lambda receives a summarize message. */ (event: "summarizeMessage", listener: (summarizeMessage: ISequencedDocumentAugmentedMessage) => void): any; /** * Emitted when the lambda receives a custom control message. */ (event: "controlMessage", listener: (controlMessage: IControlMessage) => void): any; /** * Emitted when the lambda is closing. */ (event: "close", listener: (type: LambdaCloseType) => void): any; /** * NoClient message received */ (event: "noClient", listener: () => void): any; } /** * @internal */ export declare class DeliLambda extends TypedEventEmitter implements IPartitionLambda { private readonly context; private readonly tenantId; private readonly documentId; readonly lastCheckpoint: IDeliState; private readonly clientManager; private readonly deltasProducer; private readonly signalsProducer; private readonly rawDeltasProducer; private readonly serviceConfiguration; private sessionMetric; private readonly checkpointService; private readonly sequencedSignalClients; private sequenceNumber; private signalClientConnectionNumber; private durableSequenceNumber; private logOffset; private readonly clientSeqManager; private minimumSequenceNumber; private readonly checkpointContext; private lastSendP; private lastNoClientP; private lastSentMSN; private lastHash; private lastInstruction; private lastMessageType; private activityIdleTimer; private readClientIdleTimer; private noopEvent; /** * Used for controlling op event logic */ private readonly opEvent; /** * Used for controlling checkpoint logic */ private readonly documentCheckpointManager; private globalCheckpointOnly; private readonly localCheckpointEnabled; private recievedNoClientOp; private closed; private readonly nackMessages; private serviceSummaryGenerated; constructor(context: IContext, tenantId: string, documentId: string, lastCheckpoint: IDeliState, checkpointManager: IDeliCheckpointManager, clientManager: IClientManager | undefined, deltasProducer: IProducer, signalsProducer: IProducer | undefined, rawDeltasProducer: IProducer, serviceConfiguration: IServiceConfiguration, sessionMetric: Lumber | undefined, checkpointService: ICheckpointService | undefined, sequencedSignalClients?: Map); /** * {@inheritDoc IPartitionLambda.handler} */ handler(rawMessage: IQueuedMessage): undefined; close(closeType: LambdaCloseType): void; private produceMessage; private produceMessages; private logSessionEndMetrics; private ticket; private extractDataContent; private isInvalidMessage; private createOutputMessage; private checkOrder; /** * Sends a message to the rawdeltas queue. * This essentially sends the message to this deli lambda */ private sendToRawDeltas; /** * Check if there are any old/idle write clients. * Craft and send a leave message if one is found. * To prevent recurrent leave message sending, leave messages are only piggybacked with other message type. */ private checkIdleWriteClients; /** * Check if there are any expired read clients. * The read client will expire if alfred has not sent * an ExtendClient control message within the time for 'clientTimeout'. * Craft and send a leave message for each one found. */ private checkIdleReadClients; /** * Creates a leave message for inactive clients. */ private createLeaveMessage; /** * Creates a nack message for clients. */ private createNackMessage; /** * Creates a signal message for clients. */ private createSignalMessage; private createOpMessage; private createRawOperationMessage; /** * Creates a new trace */ private createTrace; /** * Generates a checkpoint for the current state */ private generateCheckpoint; private generateDeliCheckpoint; /** * Returns a new sequence number */ private revSequenceNumber; /** * Get idle client. */ private getIdleClient; private setActivityIdleTimer; private clearActivityIdleTimer; private setReadClientIdleTimer; private clearReadClientIdleTimer; private setNoopConsolidationTimer; private clearNoopConsolidationTimer; /** * Reset the op event idle timer * Called after a message is sequenced */ private updateOpIdleTimer; private clearOpIdleTimer; /** * Resets the op event MaxTime timer * Called after an opEvent is emitted */ private updateOpMaxTimeTimer; private clearOpMaxTimeTimer; /** * Emits an opEvent for the provided type * Also resets the MaxTime timer */ private emitOpEvent; /** * Checks if the nackMessages flag should be reset */ private checkNackMessagesState; /** * Determines a checkpoint reason based on some heuristics * @returns a reason when it's time to checkpoint, or undefined if no checkpoint should be made */ private getCheckpointReason; /** * Checkpoints the current state once the pending kafka messages are produced */ private checkpoint; private readonly idleTimeCheckpoint; /** * Updates the durable sequence number * @param dsn - New durable sequence number */ private updateDurableSequenceNumber; /** * Adds/updates/removes a nack message * @param type - Nack message type * @param contents - Nack messages contents or undefined to delete the nack message */ private updateNackMessages; /** * Adds a sequenced signal client to the in-memory map. * Alfred will periodically send ExtendClient control messages, which will extend the client expiration times. * @param clientJoinMessage - Client join message (from dataContent) * @param signalMessage - Ticketed join signal message */ private addSequencedSignalClient; } //# sourceMappingURL=lambda.d.ts.map