From aecfe699f4362407a1aaadab82c58b236e3b1e2a Mon Sep 17 00:00:00 2001 From: Etienne <45695613+etiennejouan@users.noreply.github.com> Date: Wed, 13 May 2026 17:09:08 +0200 Subject: [PATCH] feat(ai-chat) - Stop ai thinking if credits exhausted (#20526) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Billing is now decremented per-step, not per-turn. The onStepFinish callback in chat-execution.service.ts calls a new decrementAndCheckAvailableCredits method on each model step, so Redis is debited incrementally as the agent runs rather than all at once at the end. Credit exhaustion stops the stream mid-run. When a step depletes the remaining credits, a hasNoMoreAvailableCredits flag is set and passed into the stopWhen predicate of streamText, causing the agent to halt before starting the next step. A new credits-exhausted event is introduced. After the stream drains and the response is persisted, if credits ran out the job publishes a dedicated credits-exhausted event to the frontend instead of the normal message-persisted event. The frontend handles this new event. useAgentChatSubscription has a new credits-exhausted case that sets a BILLING_CREDITS_EXHAUSTED-coded error on the atom, closes the writer, and stops the streaming state — triggering the existing AiChatCreditsExhaustedMessage UI. --- .../ai/hooks/useAgentChatSubscription.ts | 17 +++- .../InformationBannerEndTrialPeriod.tsx | 2 +- .../billing/services/billing-usage.service.ts | 14 ++-- .../logic-function-executor.service.ts | 2 +- .../services/agent-async-executor.service.ts | 82 ++++++++++++++++--- .../types/agent-execution-result.type.ts | 1 + .../__tests__/ai-billing.service.spec.ts | 4 +- .../ai-billing/services/ai-billing.service.ts | 39 +++++++-- .../ai/ai-chat/jobs/stream-agent-chat.job.ts | 12 ++- .../services/chat-execution.service.ts | 53 ++++++++++-- .../ai-agent/ai-agent.workflow-action.ts | 31 ++++--- ...orkflow-executor.workspace-service.spec.ts | 2 +- .../workflow-executor.workspace-service.ts | 2 +- .../ai/types/AgentChatSubscriptionEvent.ts | 3 +- 14 files changed, 212 insertions(+), 52 deletions(-) diff --git a/packages/twenty-front/src/modules/ai/hooks/useAgentChatSubscription.ts b/packages/twenty-front/src/modules/ai/hooks/useAgentChatSubscription.ts index 6836ff906c..9c97530909 100644 --- a/packages/twenty-front/src/modules/ai/hooks/useAgentChatSubscription.ts +++ b/packages/twenty-front/src/modules/ai/hooks/useAgentChatSubscription.ts @@ -19,10 +19,11 @@ import { agentChatIsStreamingComponentFamilyState } from '@/ai/states/agentChatI import { agentChatMessagesComponentFamilyState } from '@/ai/states/agentChatMessagesComponentFamilyState'; import { agentChatUsageComponentFamilyState } from '@/ai/states/agentChatUsageComponentFamilyState'; import { currentAiChatThreadTitleComponentFamilyState } from '@/ai/states/currentAiChatThreadTitleComponentFamilyState'; +import { AiChatErrorCode } from '@/ai/utils/aiChatErrorCode'; import { dispatchBrowserEvent } from '@/browser-event/utils/dispatchBrowserEvent'; +import { sseClientState } from '@/sse-db-event/states/sseClientState'; import { useAtomComponentFamilyStateCallbackState } from '@/ui/utilities/state/jotai/hooks/useAtomComponentFamilyStateCallbackState'; import { useAtomStateValue } from '@/ui/utilities/state/jotai/hooks/useAtomStateValue'; -import { sseClientState } from '@/sse-db-event/states/sseClientState'; const THROTTLE_MS = 100; @@ -313,6 +314,20 @@ export const useAgentChatSubscription = (threadId: string | null) => { store.set(isStreamingAtom, false); break; } + + case 'credits-exhausted': { + const noMoreCreditsError = new Error( + 'Chat stopped: no more available credits.', + ) as Error & { code?: string }; + + noMoreCreditsError.code = AiChatErrorCode.BILLING_CREDITS_EXHAUSTED; + store.set(errorAtom, noMoreCreditsError); + + closeWriter(); + dispatchBrowserEvent(AGENT_CHAT_REFETCH_MESSAGES_EVENT_NAME); + store.set(isStreamingAtom, false); + break; + } } }; diff --git a/packages/twenty-front/src/modules/information-banner/components/billing/InformationBannerEndTrialPeriod.tsx b/packages/twenty-front/src/modules/information-banner/components/billing/InformationBannerEndTrialPeriod.tsx index bd8c5fa31c..d30c5d4b4b 100644 --- a/packages/twenty-front/src/modules/information-banner/components/billing/InformationBannerEndTrialPeriod.tsx +++ b/packages/twenty-front/src/modules/information-banner/components/billing/InformationBannerEndTrialPeriod.tsx @@ -1,5 +1,5 @@ -import { useEndSubscriptionTrialPeriod } from '@/settings/billing/hooks/useEndSubscriptionTrialPeriod'; import { InformationBanner } from '@/information-banner/components/InformationBanner'; +import { useEndSubscriptionTrialPeriod } from '@/settings/billing/hooks/useEndSubscriptionTrialPeriod'; import { usePermissionFlagMap } from '@/settings/roles/hooks/usePermissionFlagMap'; import { useLingui } from '@lingui/react/macro'; import { PermissionFlagType } from '~/generated-metadata/graphql'; diff --git a/packages/twenty-server/src/engine/core-modules/billing/services/billing-usage.service.ts b/packages/twenty-server/src/engine/core-modules/billing/services/billing-usage.service.ts index 84736edcc0..78d6e16eb5 100644 --- a/packages/twenty-server/src/engine/core-modules/billing/services/billing-usage.service.ts +++ b/packages/twenty-server/src/engine/core-modules/billing/services/billing-usage.service.ts @@ -285,7 +285,7 @@ export class BillingUsageService { ); } - private async warmAvailableCredits( + private async warmAvailableCreditsInCache( workspaceId: string, periodStart: Date | string, periodEnd: Date | string, @@ -421,13 +421,13 @@ export class BillingUsageService { return Number(resourceCreditPrice.metadata?.credit_amount ?? 0); } - async decrementAvailableCredits({ + async decrementAvailableCreditsInCache({ workspaceId, usedCredits, }: { workspaceId: string; usedCredits: number; - }): Promise { + }): Promise { const { billingSubscription: { currentPeriodStart, currentPeriodEnd }, } = await this.workspaceCacheService.getOrRecompute(workspaceId, [ @@ -447,7 +447,7 @@ export class BillingUsageService { }); if (!isDefined(cachedAvailableCredits)) { - await this.warmAvailableCredits( + await this.warmAvailableCreditsInCache( workspaceId, currentPeriodStart, currentPeriodEnd, @@ -470,9 +470,11 @@ export class BillingUsageService { true, ); } + + return decrementedAvailableCredits; } - async invalidateAvailableCredits( + async invalidateAvailableCreditsInCache( workspaceId: string, periodStart: Date, ): Promise { @@ -504,7 +506,7 @@ export class BillingUsageService { currentPeriodStart: subscription.currentPeriodStart, }); - await this.warmAvailableCredits( + await this.warmAvailableCreditsInCache( subscription.workspaceId, subscription.currentPeriodStart, subscription.currentPeriodEnd, diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-executor/logic-function-executor.service.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-executor/logic-function-executor.service.ts index 8a14f825f6..1116a68eef 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-executor/logic-function-executor.service.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-executor/logic-function-executor.service.ts @@ -376,7 +376,7 @@ export class LogicFunctionExecutorService { periodStart = currentPeriodStart; - await this.billingUsageService.decrementAvailableCredits({ + await this.billingUsageService.decrementAvailableCreditsInCache({ workspaceId, usedCredits: 100, }); diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/services/agent-async-executor.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/services/agent-async-executor.service.ts index f0fca5e6d6..70355d1c76 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/services/agent-async-executor.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/services/agent-async-executor.service.ts @@ -24,21 +24,25 @@ import { UsageOperationType } from 'src/engine/core-modules/usage/enums/usage-op import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { WORKFLOW_AGENT_REGISTRY_TOOL_CATEGORIES } from 'src/engine/metadata-modules/ai/ai-agent-execution/constants/workflow-agent-registry-tool-categories.const'; import { type AgentExecutionResult } from 'src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type'; -import { AiBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service'; -import { countNativeWebSearchCallsFromSteps } from 'src/engine/metadata-modules/ai/ai-billing/utils/count-native-web-search-calls-from-steps.util'; -import { extractCacheCreationTokensFromSteps } from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util'; -import { mergeLanguageModelUsage } from 'src/engine/metadata-modules/ai/ai-billing/utils/merge-language-model-usage.util'; -import { - AiException, - AiExceptionCode, -} from 'src/engine/metadata-modules/ai/ai.exception'; import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-config.const'; import { WORKFLOW_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-system-prompts.const'; import { type AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity'; import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util'; +import { AiBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service'; +import { convertDollarsToBillingCredits } from 'src/engine/metadata-modules/ai/ai-billing/utils/convert-dollars-to-billing-credits.util'; +import { countNativeWebSearchCallsFromSteps } from 'src/engine/metadata-modules/ai/ai-billing/utils/count-native-web-search-calls-from-steps.util'; +import { + extractCacheCreationTokens, + extractCacheCreationTokensFromSteps, +} from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util'; +import { mergeLanguageModelUsage } from 'src/engine/metadata-modules/ai/ai-billing/utils/merge-language-model-usage.util'; import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const'; import { AiModelConfigService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-config.service'; import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service'; +import { + AiException, + AiExceptionCode, +} from 'src/engine/metadata-modules/ai/ai.exception'; import { RoleTargetEntity } from 'src/engine/metadata-modules/role-target/role-target.entity'; import { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config'; @@ -222,14 +226,35 @@ export class AgentAsyncExecutorService { this.logger.log(`Generated ${Object.keys(tools).length} tools for agent`); + let hasNoMoreAvailableCredits = false; + const textResponse = await generateText({ system: `${WORKFLOW_SYSTEM_PROMPTS.BASE}\n\n${agent ? agent.prompt : ''}`, tools, model: registeredModel.model, prompt: userPrompt, - stopWhen: stepCountIs(AGENT_CONFIG.MAX_STEPS), + stopWhen: (step) => + stepCountIs(AGENT_CONFIG.MAX_STEPS)(step) || + hasNoMoreAvailableCredits, providerOptions, experimental_telemetry: AI_TELEMETRY_CONFIG, + onStepFinish: async (step) => { + const { hasNoMoreAvailableCredits: stepHasNoMoreAvailableCredits } = + await this.aiBillingService.decrementAndCheckAvailableCredits( + registeredModel.modelId, + { + usage: step.usage, + cacheCreationTokens: extractCacheCreationTokens( + step.providerMetadata, + ), + }, + workspaceId, + ); + + if (stepHasNoMoreAvailableCredits) { + hasNoMoreAvailableCredits = true; + } + }, experimental_repairToolCall: async ({ toolCall, tools: toolsForRepair, @@ -265,6 +290,7 @@ export class AgentAsyncExecutorService { usage: textResponse.usage, cacheCreationTokens, nativeWebSearchCallCount, + hasNoMoreAvailableCredits, }; } @@ -278,6 +304,23 @@ export class AgentAsyncExecutorService { Please generate the structured output based on the execution results and context above.`, output: Output.object({ schema: jsonSchema(agentSchema) }), experimental_telemetry: AI_TELEMETRY_CONFIG, + onStepFinish: async (step) => { + const { hasNoMoreAvailableCredits: stepHasNoMoreAvailableCredits } = + await this.aiBillingService.decrementAndCheckAvailableCredits( + registeredModel.modelId, + { + usage: step.usage, + cacheCreationTokens: extractCacheCreationTokens( + step.providerMetadata, + ), + }, + workspaceId, + ); + + if (stepHasNoMoreAvailableCredits) { + hasNoMoreAvailableCredits = true; + } + }, }); accumulatedUsage = mergeLanguageModelUsage( @@ -297,6 +340,7 @@ export class AgentAsyncExecutorService { usage: accumulatedUsage, cacheCreationTokens, nativeWebSearchCallCount, + hasNoMoreAvailableCredits, }; } catch (error) { if (error instanceof AiException) { @@ -307,10 +351,24 @@ export class AgentAsyncExecutorService { AiExceptionCode.AGENT_EXECUTION_FAILED, ); } finally { - void this.aiBillingService.calculateAndBillUsage( - agent?.modelId ?? AUTO_SELECT_SMART_MODEL_ID, - { usage: accumulatedUsage, cacheCreationTokens }, + const modelId = agent?.modelId ?? AUTO_SELECT_SMART_MODEL_ID; + const costInDollars = this.aiBillingService.calculateCost(modelId, { + usage: accumulatedUsage, + cacheCreationTokens, + }); + const creditsUsedMicro = Math.round( + convertDollarsToBillingCredits(costInDollars), + ); + const totalTokens = + (accumulatedUsage.inputTokens ?? 0) + + (accumulatedUsage.outputTokens ?? 0) + + cacheCreationTokens; + + void this.aiBillingService.emitAiTokenUsageEvent( workspaceId, + creditsUsedMicro, + totalTokens, + modelId, operationType, agent?.id ?? null, userWorkspaceId, diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type.ts index ed7fe7a1a8..8b9733c8a6 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type.ts @@ -5,4 +5,5 @@ export interface AgentExecutionResult { usage: LanguageModelUsage; cacheCreationTokens: number; nativeWebSearchCallCount: number; + hasNoMoreAvailableCredits: boolean; } diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/__tests__/ai-billing.service.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/__tests__/ai-billing.service.spec.ts index da3136137e..43b47986ae 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/__tests__/ai-billing.service.spec.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/__tests__/ai-billing.service.spec.ts @@ -85,7 +85,9 @@ describe('AiBillingService', () => { { provide: BillingUsageService, useValue: { - decrementAvailableCredits: jest.fn().mockResolvedValue(undefined), + decrementAvailableCreditsInCache: jest + .fn() + .mockResolvedValue(undefined), }, }, { diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service.ts index 730768b626..ea961d7651 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service.ts @@ -76,6 +76,13 @@ export class AiBillingService { (billingInput.usage.outputTokens ?? 0) + (billingInput.cacheCreationTokens ?? 0); + if (this.billingService.isBillingEnabled()) { + await this.billingUsageService.decrementAvailableCreditsInCache({ + workspaceId, + usedCredits: creditsUsedMicro, + }); + } + await this.emitAiTokenUsageEvent( workspaceId, creditsUsedMicro, @@ -87,6 +94,29 @@ export class AiBillingService { ); } + async decrementAndCheckAvailableCredits( + modelId: ModelId, + billingInput: BillingUsageInput, + workspaceId: string, + ): Promise<{ hasNoMoreAvailableCredits: boolean }> { + if (!this.billingService.isBillingEnabled()) { + return { hasNoMoreAvailableCredits: false }; + } + + const costInDollars = this.calculateCost(modelId, billingInput); + const creditsUsedMicro = Math.round( + convertDollarsToBillingCredits(costInDollars), + ); + + const remainingCredits = + await this.billingUsageService.decrementAvailableCreditsInCache({ + workspaceId, + usedCredits: creditsUsedMicro, + }); + + return { hasNoMoreAvailableCredits: remainingCredits <= 0 }; + } + async billNativeWebSearchUsage( nativeWebSearchCallCount: number, workspaceId: string, @@ -117,7 +147,7 @@ export class AiBillingService { periodStart = currentPeriodStart; - await this.billingUsageService.decrementAvailableCredits({ + await this.billingUsageService.decrementAvailableCreditsInCache({ workspaceId, usedCredits: creditsUsedMicro, }); @@ -140,7 +170,7 @@ export class AiBillingService { ); } - private async emitAiTokenUsageEvent( + async emitAiTokenUsageEvent( workspaceId: string, creditsUsedMicro: number, totalTokens: number, @@ -159,11 +189,6 @@ export class AiBillingService { ]); periodStart = currentPeriodStart; - - await this.billingUsageService.decrementAvailableCredits({ - workspaceId, - usedCredits: creditsUsedMicro, - }); } this.workspaceEventEmitter.emitCustomBatchEvent( diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts index fd38a26034..89ab2d06f9 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts @@ -193,6 +193,7 @@ export class StreamAgentChatJob { let lastStepConversationSize = 0; let totalCacheCreationTokens = 0; let streamError: unknown; + let checkHasNoMoreAvailableCredits: () => boolean = () => false; // onFinish fires before the uiStream is fully drained. We use this // promise to coordinate: the IIFE waits for DB persist to complete @@ -224,7 +225,7 @@ export class StreamAgentChatJob { }); }; - const { stream, modelConfig } = + const { stream, modelConfig, hasNoMoreAvailableCredits } = await this.chatExecutionService.streamChat({ workspace, userWorkspaceId: data.userWorkspaceId, @@ -237,6 +238,8 @@ export class StreamAgentChatJob { conversationSizeTokens: data.conversationSizeTokens, }); + checkHasNoMoreAvailableCredits = hasNoMoreAvailableCredits; + const titleWritePromise = titlePromise.then((generatedTitle) => { if (generatedTitle) { writer.write({ @@ -315,6 +318,13 @@ export class StreamAgentChatJob { if (streamError) { reject(streamError); + } else if (checkHasNoMoreAvailableCredits()) { + await this.eventPublisherService.publish({ + threadId: data.threadId, + workspaceId: data.workspaceId, + event: { type: 'credits-exhausted' }, + }); + resolve(); } else { await this.eventPublisherService.publish({ threadId: data.threadId, diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts index e3dc901ee9..47cfba114d 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts @@ -38,8 +38,12 @@ import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/ import { type BrowsingContextType } from 'src/engine/metadata-modules/ai/ai-agent/types/browsingContext.type'; import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util'; import { AiBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service'; +import { convertDollarsToBillingCredits } from 'src/engine/metadata-modules/ai/ai-billing/utils/convert-dollars-to-billing-credits.util'; import { countNativeWebSearchCallsFromSteps } from 'src/engine/metadata-modules/ai/ai-billing/utils/count-native-web-search-calls-from-steps.util'; -import { extractCacheCreationTokensFromSteps } from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util'; +import { + extractCacheCreationTokens, + extractCacheCreationTokensFromSteps, +} from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util'; import { AI_CHAT_TOOL_NAMES_TO_PRELOAD } from 'src/engine/metadata-modules/ai/ai-chat/constants/ai-chat-tool-names-to-preload.const'; import { MessagePruningService } from 'src/engine/metadata-modules/ai/ai-chat/services/message-pruning.service'; import { SystemPromptBuilderService } from 'src/engine/metadata-modules/ai/ai-chat/services/system-prompt-builder.service'; @@ -48,9 +52,9 @@ import { type ExtractedFile, } from 'src/engine/metadata-modules/ai/ai-chat/utils/extract-code-interpreter-files.util'; import { - injectCacheBreakpoint, getCacheProviderOptions, getCallLevelCacheProviderOptions, + injectCacheBreakpoint, } from 'src/engine/metadata-modules/ai/ai-chat/utils/inject-cache-breakpoint.util'; import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const'; import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service'; @@ -72,6 +76,7 @@ export type ChatExecutionOptions = { export type ChatExecutionResult = { stream: ReturnType; modelConfig: AiModelConfig; + hasNoMoreAvailableCredits: () => boolean; }; @Injectable() @@ -262,7 +267,9 @@ export class ChatExecutionService { const modelMessages = pruningResult.messages; - const billUsageFromSteps = async (steps: StepResult[]) => { + let hasNoMoreAvailableCredits = false; + + const emitTurnUsageEvent = async (steps: StepResult[]) => { const usage = steps.reduce( (acc, step) => ({ inputTokens: (acc.inputTokens ?? 0) + (step.usage.inputTokens ?? 0), @@ -303,11 +310,24 @@ export class ChatExecutionService { ); const cacheCreationTokens = extractCacheCreationTokensFromSteps(steps); + const totalTokens = + (usage.inputTokens ?? 0) + + (usage.outputTokens ?? 0) + + cacheCreationTokens; - await this.aiBillingService.calculateAndBillUsage( + const costInDollars = this.aiBillingService.calculateCost( registeredModel.modelId, { usage, cacheCreationTokens }, + ); + const creditsUsedMicro = Math.round( + convertDollarsToBillingCredits(costInDollars), + ); + + await this.aiBillingService.emitAiTokenUsageEvent( workspace.id, + creditsUsedMicro, + totalTokens, + registeredModel.modelId, UsageOperationType.AI_CHAT_TOKEN, null, userWorkspaceId, @@ -327,7 +347,8 @@ export class ChatExecutionService { messages: [systemMessage, ...modelMessages], tools: activeTools, abortSignal, - stopWhen: stepCountIs(AGENT_CONFIG.MAX_STEPS), + stopWhen: (step) => + stepCountIs(AGENT_CONFIG.MAX_STEPS)(step) || hasNoMoreAvailableCredits, experimental_telemetry: AI_TELEMETRY_CONFIG, providerOptions: getCallLevelCacheProviderOptions( registeredModel.sdkPackage, @@ -335,8 +356,25 @@ export class ChatExecutionService { prepareStep: ({ messages }) => ({ messages: injectCacheBreakpoint(messages, registeredModel.sdkPackage), }), + onStepFinish: async (step) => { + const { hasNoMoreAvailableCredits: stepHasNoMoreAvailableCredits } = + await this.aiBillingService.decrementAndCheckAvailableCredits( + registeredModel.modelId, + { + usage: step.usage, + cacheCreationTokens: extractCacheCreationTokens( + step.providerMetadata, + ), + }, + workspace.id, + ); + + if (stepHasNoMoreAvailableCredits) { + hasNoMoreAvailableCredits = true; + } + }, onAbort: async ({ steps }) => { - await billUsageFromSteps(steps); + await emitTurnUsageEvent(steps); }, experimental_repairToolCall: async ({ toolCall, @@ -363,7 +401,7 @@ export class ChatExecutionService { Promise.all([stream.usage, stream.steps]) .then(async ([, steps]) => { - await billUsageFromSteps(steps); + await emitTurnUsageEvent(steps); }) .catch((error) => { if (error?.name === 'AbortError') { @@ -375,6 +413,7 @@ export class ChatExecutionService { return { stream, modelConfig, + hasNoMoreAvailableCredits: () => hasNoMoreAvailableCredits, }; } diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts index 5676037004..33f0e3879f 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts @@ -76,18 +76,25 @@ export class AiAgentWorkflowAction implements WorkflowAction { ? executionContext.authContext.userWorkspaceId : null; - const { result } = await this.aiAgentExecutionService.executeAgent({ - agent, - userPrompt: resolveInput(prompt, context) as string, - actorContext: executionContext.isActingOnBehalfOfUser - ? executionContext.initiator - : undefined, - rolePermissionConfig: executionContext.rolePermissionConfig, - authContext: executionContext.authContext, - workspaceId, - userWorkspaceId, - operationType: UsageOperationType.AI_WORKFLOW_TOKEN, - }); + const { result, hasNoMoreAvailableCredits } = + await this.aiAgentExecutionService.executeAgent({ + agent, + userPrompt: resolveInput(prompt, context) as string, + actorContext: executionContext.isActingOnBehalfOfUser + ? executionContext.initiator + : undefined, + rolePermissionConfig: executionContext.rolePermissionConfig, + authContext: executionContext.authContext, + workspaceId, + userWorkspaceId, + operationType: UsageOperationType.AI_WORKFLOW_TOKEN, + }); + + if (hasNoMoreAvailableCredits) { + return { + error: 'AI agent stopped: no more available credits.', + }; + } return { result, diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts index 424b95119e..b30eb65550 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts @@ -78,7 +78,7 @@ describe('WorkflowExecutorWorkspaceService', () => { const mockBillingUsageService = { hasAvailableCredits: jest.fn().mockResolvedValue(true), - decrementAvailableCredits: jest.fn().mockResolvedValue(undefined), + decrementAvailableCreditsInCache: jest.fn().mockResolvedValue(undefined), }; const mockExceptionHandlerService = { diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts index 02bee70806..bda731f7ec 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts @@ -375,7 +375,7 @@ export class WorkflowExecutorWorkspaceService { periodStart = currentPeriodStart; - await this.billingUsageService.decrementAvailableCredits({ + await this.billingUsageService.decrementAvailableCreditsInCache({ workspaceId, usedCredits: 100, }); diff --git a/packages/twenty-shared/src/ai/types/AgentChatSubscriptionEvent.ts b/packages/twenty-shared/src/ai/types/AgentChatSubscriptionEvent.ts index 8c2343ee7b..7ec44a20b5 100644 --- a/packages/twenty-shared/src/ai/types/AgentChatSubscriptionEvent.ts +++ b/packages/twenty-shared/src/ai/types/AgentChatSubscriptionEvent.ts @@ -5,4 +5,5 @@ export type AgentChatSubscriptionEvent = | { type: 'stream-chunk'; chunk: Record; seq?: number } | { type: 'message-persisted'; messageId: string } | { type: 'queue-updated' } - | { type: 'stream-error'; code: string; message: string }; + | { type: 'stream-error'; code: string; message: string } + | { type: 'credits-exhausted' };