From 35efbcc9823116072a89cc3c7ca371000f505e20 Mon Sep 17 00:00:00 2001 From: Ravi Kumar L Date: Wed, 29 Jul 2026 20:20:50 +0200 Subject: [PATCH] fix(langfuse): preserve trace sampling for feedback --- api/app/clients/BaseClient.js | 5 ++++ api/app/clients/specs/BaseClient.test.js | 24 +++++++++++++++++ api/server/routes/messages.js | 1 + packages/api/src/langfuse/destinations.ts | 8 +++++- packages/api/src/langfuse/feedback.spec.ts | 26 +++++++++++++++++++ packages/api/src/langfuse/feedback.ts | 4 ++- .../data-schemas/src/methods/message.spec.ts | 14 ++++++++++ packages/data-schemas/src/methods/message.ts | 1 + packages/data-schemas/src/schema/message.ts | 3 +++ packages/data-schemas/src/types/message.ts | 1 + 10 files changed, 85 insertions(+), 2 deletions(-) diff --git a/api/app/clients/BaseClient.js b/api/app/clients/BaseClient.js index 6d541e160f..214c4abd00 100644 --- a/api/app/clients/BaseClient.js +++ b/api/app/clients/BaseClient.js @@ -12,6 +12,8 @@ const { encodeAndFormatAudios, encodeAndFormatVideos, encodeAndFormatDocuments, + isLangfuseTraceSampled, + traceIdForMessage, } = require('@librechat/api'); const { Constants, @@ -716,6 +718,9 @@ class BaseClient { conversationId, parentMessageId: userMessage.messageId, isCreatedByUser: false, + ...(isAgentsEndpoint(this.options.endpoint) && { + langfuseSampled: isLangfuseTraceSampled(traceIdForMessage(responseMessageId)), + }), isEdited, model: this.getResponseModel(), sender: this.sender, diff --git a/api/app/clients/specs/BaseClient.test.js b/api/app/clients/specs/BaseClient.test.js index d565f87012..8c22181599 100644 --- a/api/app/clients/specs/BaseClient.test.js +++ b/api/app/clients/specs/BaseClient.test.js @@ -746,6 +746,30 @@ describe('BaseClient', () => { ); }); + test('persists the generation-time Langfuse sampling decision for agent responses', async () => { + const previousSampleRate = process.env.LANGFUSE_SAMPLE_RATE; + process.env.LANGFUSE_SAMPLE_RATE = '0'; + TestClient.options.endpoint = 'agents'; + const saveSpy = jest.spyOn(TestClient, 'saveMessageToDatabase'); + + try { + const response = await TestClient.sendMessage('Hello, world!', { user: {} }); + + expect(response.langfuseSampled).toBe(false); + expect(saveSpy).toHaveBeenCalledWith( + expect.objectContaining({ langfuseSampled: false }), + expect.any(Object), + expect.any(Object), + ); + } finally { + if (previousSampleRate == null) { + delete process.env.LANGFUSE_SAMPLE_RATE; + } else { + process.env.LANGFUSE_SAMPLE_RATE = previousSampleRate; + } + } + }); + test('should handle existing conversation when getConvo retrieves one', async () => { const existingConvo = { conversationId: 'existing-convo-id', diff --git a/api/server/routes/messages.js b/api/server/routes/messages.js index f541202989..b52be480a1 100644 --- a/api/server/routes/messages.js +++ b/api/server/routes/messages.js @@ -445,6 +445,7 @@ router.put( if (!isAssistantsEndpoint(updatedMessage.endpoint)) { sendFeedbackScore({ traceId: traceIdForMessage(messageId), + sampled: updatedMessage.langfuseSampled, feedback: updatedMessage.feedback, appConfig: req.config, metadata: { diff --git a/packages/api/src/langfuse/destinations.ts b/packages/api/src/langfuse/destinations.ts index 3a36683416..c128049cd2 100644 --- a/packages/api/src/langfuse/destinations.ts +++ b/packages/api/src/langfuse/destinations.ts @@ -3,6 +3,7 @@ import { hasLangfuseEnvCredentials, isLangfuseFanoutEnabled, isLangfuseTenantExportEnabled, + isLangfuseTracingEnabled, isLangfuseTraceSampled, usesLangfuseMultiTenantRouting, } from './policy'; @@ -106,8 +107,13 @@ function getConfiguredScoreDestination( export function getScoreDestinations( appConfig: AppConfig | undefined, traceId: string, + sampled?: boolean, ): LangfuseScoreDestination[] { - if (!isLangfuseTraceSampled(traceId)) { + if ( + !isLangfuseTracingEnabled() || + sampled === false || + (sampled == null && !isLangfuseTraceSampled(traceId)) + ) { return []; } diff --git a/packages/api/src/langfuse/feedback.spec.ts b/packages/api/src/langfuse/feedback.spec.ts index d3c96ff01f..babfe99544 100644 --- a/packages/api/src/langfuse/feedback.spec.ts +++ b/packages/api/src/langfuse/feedback.spec.ts @@ -256,6 +256,32 @@ describe('Langfuse feedback scores', () => { expect(getFetchMock()).not.toHaveBeenCalled(); }); + it('preserves a sampled trace when the sample rate decreases', async () => { + process.env.LANGFUSE_SAMPLE_RATE = '0.1'; + const { sendFeedbackScore } = await loadFeedback(); + + await sendFeedbackScore({ + traceId: '658f74b0a232417fc3e6e4d9ef5f563a', + sampled: true, + feedback: { rating: 'thumbsUp' }, + }); + + expect(getFetchMock()).toHaveBeenCalledTimes(1); + }); + + it('preserves an excluded trace when the sample rate increases', async () => { + process.env.LANGFUSE_SAMPLE_RATE = '1'; + const { sendFeedbackScore } = await loadFeedback(); + + await sendFeedbackScore({ + traceId: '86d413435f8b0d7f32d4d010ce769e2e', + sampled: false, + feedback: { rating: 'thumbsUp' }, + }); + + expect(getFetchMock()).not.toHaveBeenCalled(); + }); + it('posts feedback scores to central fanout and tenant Langfuse projects', async () => { enableTenantFanout(); delete process.env.TENANT_ISOLATION_STRICT; diff --git a/packages/api/src/langfuse/feedback.ts b/packages/api/src/langfuse/feedback.ts index a355e34007..8355437ff7 100644 --- a/packages/api/src/langfuse/feedback.ts +++ b/packages/api/src/langfuse/feedback.ts @@ -12,6 +12,7 @@ export type LangfuseFeedbackMetadata = Record { expect(updatedMessage?.text).toBe('Updated text'); }); + it('returns the generation-time Langfuse sampling decision with feedback updates', async () => { + await saveMessage(mockCtx, { + ...mockMessageData, + langfuseSampled: true, + }); + + const result = await updateMessage(mockCtx.userId, { + messageId: 'msg123', + feedback: { rating: 'thumbsUp' }, + }); + + expect(result?.langfuseSampled).toBe(true); + }); + it('should throw an error if message is not found', async () => { await expect( updateMessage(mockCtx.userId, { messageId: 'nonexistent', text: 'Test' }), diff --git a/packages/data-schemas/src/methods/message.ts b/packages/data-schemas/src/methods/message.ts index a66666f5c6..859a12ab2c 100644 --- a/packages/data-schemas/src/methods/message.ts +++ b/packages/data-schemas/src/methods/message.ts @@ -295,6 +295,7 @@ export function createMessageMethods(mongoose: typeof import('mongoose')): Messa tokenCount: updatedMessage.tokenCount, feedback: updatedMessage.feedback, endpoint: updatedMessage.endpoint, + langfuseSampled: updatedMessage.langfuseSampled, }; } catch (err) { logger.error('Error updating message:', err); diff --git a/packages/data-schemas/src/schema/message.ts b/packages/data-schemas/src/schema/message.ts index d7fd3fafb0..84364fe71b 100644 --- a/packages/data-schemas/src/schema/message.ts +++ b/packages/data-schemas/src/schema/message.ts @@ -97,6 +97,9 @@ const messageSchema: Schema = new Schema( default: undefined, required: false, }, + langfuseSampled: { + type: Boolean, + }, _meiliIndex: { type: Boolean, required: false, diff --git a/packages/data-schemas/src/types/message.ts b/packages/data-schemas/src/types/message.ts index 96d9c1ff35..0fd2f50443 100644 --- a/packages/data-schemas/src/types/message.ts +++ b/packages/data-schemas/src/types/message.ts @@ -27,6 +27,7 @@ export interface IMessage extends Document { tag: TFeedbackTag | undefined; text?: string; }; + langfuseSampled?: boolean; _meiliIndex?: boolean; files?: unknown[]; plugin?: {