/*! * Copyright (c) Microsoft Corporation and contributors. All rights reserved. * Licensed under the MIT License. */ import { ProtocolOpHandler } from "@fluidframework/protocol-base"; import { ISequencedDocumentMessage } from "@fluidframework/protocol-definitions"; import { IContext, IProducer, IScribe, IServiceConfiguration, IQueuedMessage, IPartitionLambda, LambdaCloseType } from "@fluidframework/server-services-core"; import { Lumber, LumberEventName } from "@fluidframework/server-services-telemetry"; import { CheckpointReason } from "../utils"; import { ICheckpointManager, IPendingMessageReader, ISummaryWriter } from "./interfaces"; /** * @internal */ export declare class ScribeLambda implements IPartitionLambda { protected readonly context: IContext; protected tenantId: string; protected documentId: string; private readonly summaryWriter; private readonly pendingMessageReader; private readonly checkpointManager; private readonly serviceConfiguration; private readonly producer; private protocolHandler; private protocolHead; private scribeSessionMetric; private readonly transientTenants; private readonly disableTransientTenantFiltering; private readonly restartOnCheckpointFailure; private readonly kafkaCheckpointOnReprocessingOp; private readonly isEphemeralContainer; private readonly localCheckpointEnabled; private readonly maxPendingCheckpointMessagesLength; private lastOffset; private pendingCheckpointScribe; private pendingCheckpointOffset; private pendingP; private readonly pendingCheckpointMessages; private pendingMessages; private sequenceNumber; private minSequenceNumber; private lastClientSummaryHead; private lastSummarySequenceNumber; private validParentSummaries; private isDocumentCorrupt; private clearCache; private closed; private readonly documentCheckpointManager; private globalCheckpointOnly; private lastCheckpointInsertedNumber; constructor(context: IContext, tenantId: string, documentId: string, summaryWriter: ISummaryWriter, pendingMessageReader: IPendingMessageReader | undefined, checkpointManager: ICheckpointManager, scribe: IScribe, serviceConfiguration: IServiceConfiguration, producer: IProducer | undefined, protocolHandler: ProtocolOpHandler, protocolHead: number, messages: ISequencedDocumentMessage[], scribeSessionMetric: Lumber | undefined, transientTenants: Set, disableTransientTenantFiltering: boolean, restartOnCheckpointFailure: boolean, kafkaCheckpointOnReprocessingOp: boolean, isEphemeralContainer: boolean, localCheckpointEnabled: boolean, maxPendingCheckpointMessagesLength: number); /** * {@inheritDoc IPartitionLambda.handler} */ handler(message: IQueuedMessage): Promise; prepareCheckpoint(message: IQueuedMessage, checkpointReason: CheckpointReason, skipKafkaCheckpoint?: boolean): void; close(closeType: LambdaCloseType): void; private logScribeSessionMetrics; private processFromPending; private markDocumentAsCorrupt; private revertProtocolState; private generateScribeCheckpoint; private checkpointCore; private writeCheckpoint; /** * Protocol head is the sequence number of the last summary * This method updates the protocol head to the new summary sequence number * @param protocolHead - The sequence number of the new summary */ private updateProtocolHead; /** * lastSummarySequenceNumber tracks the sequence number that was part of the latest summary * This method updates it to the sequence number that was part of the latest summary * @param summarySequenceNumber - The sequence number of the operation that was part of the latest summary */ private updateLastSummarySequenceNumber; /** * validParentSummaries tracks summary handles for service summaries that have been written since the latest client summary. * @param summaryHandle - The handle for a service summary that occurred after latest client summary. */ private updateValidParentSummaries; private sendSummaryAck; private sendSummaryNack; private sendSummaryConfirmationMessage; private setStateFromCheckpoint; private getCheckpointReason; private readonly idleTimeCheckpoint; } //# sourceMappingURL=lambda.d.ts.map