diff --git a/front/lib/analytics/agent_message_consumption/index.ts b/front/lib/analytics/agent_message_consumption/index.ts new file mode 100644 index 000000000000..75d991d3a603 --- /dev/null +++ b/front/lib/analytics/agent_message_consumption/index.ts @@ -0,0 +1,33 @@ +import { buildAgentMessageConsumptionAnalyticsDocuments } from "@app/lib/analytics/agent_message_consumption/documents"; +import { loadAgentMessageConsumptionAnalyticsInput } from "@app/lib/analytics/agent_message_consumption/load"; +import { upsertAgentMessageConsumptionAnalyticsDocuments } from "@app/lib/analytics/agent_message_consumption/store"; +import type { ElasticsearchError } from "@app/lib/api/elasticsearch"; +import type { Authenticator } from "@app/lib/auth"; +import type { Result } from "@app/types/shared/result"; +import { Ok } from "@app/types/shared/result"; +import assert from "assert"; + +/** + * Loads, projects, and indexes the complete consumption analytics snapshot for one agent message. + * Callers only identify the message. This module owns the ordering and completeness requirements + * of the indexed snapshot. + */ +export async function indexAgentMessageConsumptionAnalytics( + auth: Authenticator, + { agentMessageId }: { agentMessageId: string } +): Promise> { + const input = await loadAgentMessageConsumptionAnalyticsInput(auth, { + agentMessageId, + }); + if (!input) { + return new Ok(undefined); + } + + const documents = buildAgentMessageConsumptionAnalyticsDocuments(input); + assert( + documents && documents.length > 0, + "Consumption attribution is incomplete for analytics" + ); + + return upsertAgentMessageConsumptionAnalyticsDocuments(documents); +} diff --git a/front/lib/analytics/agent_message_consumption/store.test.ts b/front/lib/analytics/agent_message_consumption/store.test.ts new file mode 100644 index 000000000000..dbd2474246ec --- /dev/null +++ b/front/lib/analytics/agent_message_consumption/store.test.ts @@ -0,0 +1,164 @@ +import { upsertAgentMessageConsumptionAnalyticsDocuments } from "@app/lib/analytics/agent_message_consumption/store"; +import { + CONSUMPTION_ANALYTICS_ALIAS_NAME, + ElasticsearchError, + withEs, +} from "@app/lib/api/elasticsearch"; +import { USAGE_TYPE_USER } from "@app/lib/metronome/constants"; +import type { AgentMessageConsumptionAnalyticsData } from "@app/types/assistant/analytics"; +import { Err, Ok } from "@app/types/shared/result"; +import { normalizeError } from "@app/types/shared/utils/error_utils"; +import { Client } from "@elastic/elasticsearch"; +import { afterAll, beforeEach, describe, expect, it, vi } from "vitest"; + +vi.mock("@app/lib/api/elasticsearch", async (importActual) => { + const actual = + await importActual(); + return { ...actual, withEs: vi.fn() }; +}); + +const client = new Client({ node: "http://localhost:9200" }); +const bulkMock = vi.spyOn(client, "bulk"); + +function makeDocument(): AgentMessageConsumptionAnalyticsData { + return { + agent: { + id: "agent_1", + version: "1", + tag_ids: [], + parent_ids: [], + direct_parent_id: null, + root_id: "agent_1", + depth: 0, + }, + agent_message_id: "agent_message_1", + api_key_name: null, + attribution_version: 4, + completed_at: "2026-08-07T12:00:00.000Z", + consumption_key: "run-usage:1", + consumption_type: "llm", + context_origin: "web", + conversation_id: "conversation_1", + credit_micro: 1_000_000, + execution_time_ms: null, + gross_credit_micro: { + system: 0, + input: 600_000, + result_footprint: null, + output: 400_000, + reasoning: 0, + direct: 0, + total: 1_000_000, + }, + message_version: "2", + model: null, + run_usage_id: "1", + space_id: null, + status: "succeeded", + step_index: 0, + tokens: { + system: 0, + input: 10, + result_footprint: null, + output: 5, + reasoning: 0, + }, + tool: null, + trigger_id: null, + usage_type: USAGE_TYPE_USER, + user: null, + workspace_id: "workspace_1", + }; +} + +describe("upsertAgentMessageConsumptionAnalyticsDocuments", () => { + beforeEach(() => { + vi.clearAllMocks(); + bulkMock.mockResolvedValue({ errors: false, items: [], took: 1 }); + vi.mocked(withEs).mockImplementation(async (fn) => { + try { + return new Ok(await fn(client)); + } catch (error) { + return new Err( + new ElasticsearchError("query_error", normalizeError(error).message) + ); + } + }); + }); + + afterAll(async () => { + await client.close(); + }); + + it("uses a stable identity for idempotent upserts", async () => { + const document = makeDocument(); + + const result = await upsertAgentMessageConsumptionAnalyticsDocuments([ + document, + ]); + + expect(result.isOk()).toBe(true); + expect(bulkMock).toHaveBeenCalledWith({ + body: [ + { + index: { + _index: CONSUMPTION_ANALYTICS_ALIAS_NAME, + _id: "workspace_1_agent_message_1_run-usage:1", + }, + }, + document, + ], + refresh: false, + }); + }); + + it("does nothing when there are no documents", async () => { + const result = await upsertAgentMessageConsumptionAnalyticsDocuments([]); + + expect(result.isOk()).toBe(true); + expect(withEs).not.toHaveBeenCalled(); + }); + + it("returns the Elasticsearch error when the request fails", async () => { + const error = new ElasticsearchError("connection_error", "write failed"); + vi.mocked(withEs).mockResolvedValueOnce(new Err(error)); + + const result = await upsertAgentMessageConsumptionAnalyticsDocuments([ + makeDocument(), + ]); + + expect(result.isErr()).toBe(true); + if (result.isErr()) { + expect(result.error).toBe(error); + } + }); + + it("returns the error from a failed bulk item", async () => { + bulkMock.mockResolvedValueOnce({ + errors: true, + items: [ + { + index: { + _index: CONSUMPTION_ANALYTICS_ALIAS_NAME, + status: 429, + error: { + type: "es_rejected_execution_exception", + reason: "queue full", + }, + }, + }, + ], + took: 1, + }); + + const result = await upsertAgentMessageConsumptionAnalyticsDocuments([ + makeDocument(), + ]); + + expect(result.isErr()).toBe(true); + if (result.isErr()) { + expect(result.error.message).toBe("queue full"); + expect(result.error.statusCode).toBe(429); + } + }); +}); diff --git a/front/lib/analytics/agent_message_consumption/store.ts b/front/lib/analytics/agent_message_consumption/store.ts new file mode 100644 index 000000000000..08292dc0d05c --- /dev/null +++ b/front/lib/analytics/agent_message_consumption/store.ts @@ -0,0 +1,62 @@ +import { + CONSUMPTION_ANALYTICS_ALIAS_NAME, + ElasticsearchError, + withEs, +} from "@app/lib/api/elasticsearch"; +import type { AgentMessageConsumptionAnalyticsData } from "@app/types/assistant/analytics"; +import type { Result } from "@app/types/shared/result"; +import { Err, Ok } from "@app/types/shared/result"; + +function makeAgentMessageConsumptionAnalyticsDocumentId( + document: Pick< + AgentMessageConsumptionAnalyticsData, + "agent_message_id" | "consumption_key" | "workspace_id" + > +): string { + return `${document.workspace_id}_${document.agent_message_id}_${document.consumption_key}`; +} + +/** Upserts every consumption unit using its stable identity. */ +export async function upsertAgentMessageConsumptionAnalyticsDocuments( + documents: AgentMessageConsumptionAnalyticsData[] +): Promise> { + if (documents.length === 0) { + return new Ok(undefined); + } + + const result = await withEs((client) => + client.bulk({ + body: documents.flatMap((document) => [ + { + index: { + _index: CONSUMPTION_ANALYTICS_ALIAS_NAME, + _id: makeAgentMessageConsumptionAnalyticsDocumentId(document), + }, + }, + document, + ]), + refresh: false, + }) + ); + + if (result.isErr()) { + return result; + } + + if (!result.value.errors) { + return new Ok(undefined); + } + + const failedItem = result.value.items.find( + (item) => item.index?.error + )?.index; + + return new Err( + new ElasticsearchError( + "query_error", + failedItem?.error?.reason ?? + "Elasticsearch bulk response contains failed items", + failedItem?.status + ) + ); +} diff --git a/front/temporal/agent_loop/activities/finalize.ts b/front/temporal/agent_loop/activities/finalize.ts index 8a38875bcff8..0c1b4f7ea516 100644 --- a/front/temporal/agent_loop/activities/finalize.ts +++ b/front/temporal/agent_loop/activities/finalize.ts @@ -29,6 +29,23 @@ import { } from "@app/temporal/agent_loop/activities/usage_tracking"; import type { AgentLoopArgs } from "@app/types/assistant/agent_run"; +async function launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth: Authenticator, + agentLoopArgs: AgentLoopArgs, + { + creditArgs = agentLoopArgs, + }: { + creditArgs?: { agentMessageId: string; dustRunIds?: string[] }; + } = {} +): Promise { + // Consumption analytics needs the authoritative bill, usage type, and historical skill snapshot + // before its attribution workflow can safely materialize Elasticsearch documents. + await snapshotAgentMessageSkills(auth, agentLoopArgs); + await computeAndStoreAgentMessageCredits(auth, creditArgs); + + await launchAgentMessageConsumptionAttribution(auth, agentLoopArgs); +} + export async function finalizeSuccessfulAgentLoopActivity( authType: AuthenticatorType, agentLoopArgs: AgentLoopArgs @@ -36,12 +53,13 @@ export async function finalizeSuccessfulAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, agentLoopArgs), conversationUnreadNotification(auth, agentLoopArgs), activationNewConversationNotification(auth, agentLoopArgs), handleMentions(auth, agentLoopArgs), @@ -64,12 +82,13 @@ export async function finalizeGracefullyStoppedAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, agentLoopArgs), conversationUnreadNotification(auth, agentLoopArgs), handleMentions(auth, agentLoopArgs), ]); @@ -92,12 +111,13 @@ export async function finalizeInterruptedAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, agentLoopArgs), conversationUnreadNotification(auth, agentLoopArgs), handleMentions(auth, agentLoopArgs), ]); @@ -112,12 +132,13 @@ export async function finalizeCancelledAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, agentLoopArgs), sendEmailReplyOnError( auth, agentLoopArgs, @@ -135,14 +156,16 @@ export async function finalizeCreditStoppedAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs, + { + creditArgs: { agentMessageId: agentLoopArgs.agentMessageId }, + } + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, { - agentMessageId: agentLoopArgs.agentMessageId, - }), sendEmailReplyOnError(auth, agentLoopArgs, creditsExhaustedMessage(auth)), ]); } @@ -157,12 +180,13 @@ export async function finalizeErroredAgentLoopActivity( const auth = await Authenticator.fromJsonWithRefrehedGroups(authType); await Promise.all([ - snapshotAgentMessageSkills(auth, agentLoopArgs), launchAgentMessageAnalytics(auth, agentLoopArgs), - launchAgentMessageConsumptionAttribution(auth, agentLoopArgs), + launchAgentMessageConsumptionAttributionAfterPersistingInputs( + auth, + agentLoopArgs + ), launchTrackProgrammaticUsage(auth, agentLoopArgs), launchEmitMetronomeUsageEvents(auth, agentLoopArgs), - computeAndStoreAgentMessageCredits(auth, agentLoopArgs), sendEmailReplyOnError( auth, agentLoopArgs, diff --git a/front/temporal/analytics_queue/activities/consumption_attribution.test.ts b/front/temporal/analytics_queue/activities/consumption_attribution.test.ts new file mode 100644 index 000000000000..4501098e4e93 --- /dev/null +++ b/front/temporal/analytics_queue/activities/consumption_attribution.test.ts @@ -0,0 +1,58 @@ +import { indexAgentMessageConsumptionAnalytics } from "@app/lib/analytics/agent_message_consumption"; +import { ElasticsearchError } from "@app/lib/api/elasticsearch"; +import type { AuthenticatorType } from "@app/lib/auth"; +import { Authenticator } from "@app/lib/auth"; +import { storeAgentMessageConsumptionAnalyticsActivity } from "@app/temporal/analytics_queue/activities/consumption_attribution"; +import type { AgentLoopArgs } from "@app/types/assistant/agent_run"; +import { Err, Ok } from "@app/types/shared/result"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +vi.mock( + "@app/lib/analytics/agent_message_consumption", + async (importActual) => { + const actual = + await importActual< + typeof import("@app/lib/analytics/agent_message_consumption") + >(); + return { ...actual, indexAgentMessageConsumptionAnalytics: vi.fn() }; + } +); + +const authType = {} as AuthenticatorType; +const agentLoopArgs = { + agentMessageId: "agent_message_1", +} as AgentLoopArgs; + +describe("storeAgentMessageConsumptionAnalyticsActivity", () => { + beforeEach(() => { + vi.clearAllMocks(); + vi.spyOn(Authenticator, "fromJSON").mockResolvedValue({ + getNonNullableWorkspace: () => ({ sId: "workspace_1" }), + } as Authenticator); + }); + + it("completes when indexing succeeds", async () => { + vi.mocked(indexAgentMessageConsumptionAnalytics).mockResolvedValue( + new Ok(undefined) + ); + + await expect( + storeAgentMessageConsumptionAnalyticsActivity(authType, { + agentLoopArgs, + }) + ).resolves.toBeUndefined(); + }); + + it("throws the Elasticsearch error so Temporal retries the activity", async () => { + const error = new ElasticsearchError("query_error", "invalid mapping", 400); + vi.mocked(indexAgentMessageConsumptionAnalytics).mockResolvedValue( + new Err(error) + ); + + await expect( + storeAgentMessageConsumptionAnalyticsActivity(authType, { + agentLoopArgs, + }) + ).rejects.toBe(error); + }); +}); diff --git a/front/temporal/analytics_queue/activities/consumption_attribution.ts b/front/temporal/analytics_queue/activities/consumption_attribution.ts index c12fa6cd183e..96555051e3f7 100644 --- a/front/temporal/analytics_queue/activities/consumption_attribution.ts +++ b/front/temporal/analytics_queue/activities/consumption_attribution.ts @@ -1,6 +1,8 @@ +import { indexAgentMessageConsumptionAnalytics } from "@app/lib/analytics/agent_message_consumption"; import { computeAndStoreAgentMessageConsumptionAttribution } from "@app/lib/api/assistant/agent_message_consumption_attribution/store"; import type { AuthenticatorType } from "@app/lib/auth"; import { Authenticator } from "@app/lib/auth"; +import logger from "@app/logger/logger"; import type { AgentLoopArgs } from "@app/types/assistant/agent_run"; export async function storeAgentMessageConsumptionAttributionActivity( @@ -19,3 +21,34 @@ export async function storeAgentMessageConsumptionAttributionActivity( conversationId, }); } + +/** Builds and bulk-upserts every billed consumption unit for one settled agent message. */ +export async function storeAgentMessageConsumptionAnalyticsActivity( + authType: AuthenticatorType, + { + agentLoopArgs, + }: { + agentLoopArgs: AgentLoopArgs; + } +): Promise { + const auth = await Authenticator.fromJSON(authType); + const result = await indexAgentMessageConsumptionAnalytics(auth, { + agentMessageId: agentLoopArgs.agentMessageId, + }); + + if (result.isErr()) { + const { error } = result; + const workspaceId = auth.getNonNullableWorkspace().sId; + + logger.error( + { + error, + workspaceId, + agentMessageId: agentLoopArgs.agentMessageId, + }, + "[ConsumptionAnalytics] Failed to upsert consumption documents in ES" + ); + + throw error; + } +} diff --git a/front/temporal/analytics_queue/client.test.ts b/front/temporal/analytics_queue/client.test.ts index 15e423295007..88bf4f571d44 100644 --- a/front/temporal/analytics_queue/client.test.ts +++ b/front/temporal/analytics_queue/client.test.ts @@ -1,8 +1,8 @@ import { launchStoreAgentMessageConsumptionAttributionWorkflow } from "@app/temporal/analytics_queue/client"; import { QUEUE_NAME } from "@app/temporal/analytics_queue/config"; import { makeAgentMessageAnalyticsWorkflowId } from "@app/temporal/analytics_queue/helpers"; -import { storeAgentMessageConsumptionAttributionV2Signal } from "@app/temporal/analytics_queue/signals"; -import { storeAgentMessageConsumptionAttributionV2Workflow } from "@app/temporal/analytics_queue/workflows"; +import { storeAgentMessageConsumptionAttributionV3Signal } from "@app/temporal/analytics_queue/signals"; +import { storeAgentMessageConsumptionAttributionV3Workflow } from "@app/temporal/analytics_queue/workflows"; import { createResourceTest } from "@app/tests/utils/generic_resource_tests"; import type { AgentLoopArgs } from "@app/types/assistant/agent_run"; import { beforeEach, describe, expect, it, vi } from "vitest"; @@ -25,7 +25,7 @@ describe("launchStoreAgentMessageConsumptionAttributionWorkflow", () => { mockSignalWithStart.mockResolvedValue(undefined); }); - it("signals the replay-safe V2 workflow for every finalize", async () => { + it("signals the replay-safe V3 workflow for every finalize", async () => { const { authenticator } = await createResourceTest({}); const authType = authenticator.toJSON(); const agentLoopArgs: AgentLoopArgs = { @@ -51,7 +51,7 @@ describe("launchStoreAgentMessageConsumptionAttributionWorkflow", () => { expect(second.isOk()).toBe(true); expect(mockSignalWithStart).toHaveBeenCalledTimes(2); expect(mockSignalWithStart).toHaveBeenCalledWith( - storeAgentMessageConsumptionAttributionV2Workflow, + storeAgentMessageConsumptionAttributionV3Workflow, expect.objectContaining({ args: [authType, { agentLoopArgs }], taskQueue: QUEUE_NAME, @@ -59,8 +59,8 @@ describe("launchStoreAgentMessageConsumptionAttributionWorkflow", () => { agentMessageId: agentLoopArgs.agentMessageId, conversationId: agentLoopArgs.conversationId, workspaceId: authType.workspaceId, - })}-consumption-attribution-v2`, - signal: storeAgentMessageConsumptionAttributionV2Signal, + })}-consumption-attribution-v3`, + signal: storeAgentMessageConsumptionAttributionV3Signal, signalArgs: undefined, }) ); diff --git a/front/temporal/analytics_queue/client.ts b/front/temporal/analytics_queue/client.ts index d40acc831474..ae2b975f98eb 100644 --- a/front/temporal/analytics_queue/client.ts +++ b/front/temporal/analytics_queue/client.ts @@ -7,10 +7,10 @@ import { getTemporalClientForFrontNamespace } from "@app/lib/temporal"; import logger from "@app/logger/logger"; import { QUEUE_NAME } from "@app/temporal/analytics_queue/config"; import { makeAgentMessageAnalyticsWorkflowId } from "@app/temporal/analytics_queue/helpers"; -import { storeAgentMessageConsumptionAttributionV2Signal } from "@app/temporal/analytics_queue/signals"; +import { storeAgentMessageConsumptionAttributionV3Signal } from "@app/temporal/analytics_queue/signals"; import { storeAgentAnalyticsWorkflow, - storeAgentMessageConsumptionAttributionV2Workflow, + storeAgentMessageConsumptionAttributionV3Workflow, storeAgentMessageFeedbackWorkflow, } from "@app/temporal/analytics_queue/workflows"; import type { @@ -117,7 +117,7 @@ export async function launchStoreAgentMessageConsumptionAttributionWorkflow({ agentMessageId, conversationId, workspaceId, - }) + "-consumption-attribution-v2"; + }) + "-consumption-attribution-v3"; try { // signalWithStart, not start: a message settles across several finalizes (pause for approval, @@ -125,12 +125,12 @@ export async function launchStoreAgentMessageConsumptionAttributionWorkflow({ // already-started, freezing a tool that was still blocked when the first pass ran. The signal // instead reruns a workflow already in flight and starts one otherwise. await client.workflow.signalWithStart( - storeAgentMessageConsumptionAttributionV2Workflow, + storeAgentMessageConsumptionAttributionV3Workflow, { args: [authType, { agentLoopArgs }], taskQueue: QUEUE_NAME, workflowId, - signal: storeAgentMessageConsumptionAttributionV2Signal, + signal: storeAgentMessageConsumptionAttributionV3Signal, signalArgs: undefined, searchAttributes: { conversationId: [conversationId], diff --git a/front/temporal/analytics_queue/signals.ts b/front/temporal/analytics_queue/signals.ts index a61460cfd38b..3b7ea4e2c70a 100644 --- a/front/temporal/analytics_queue/signals.ts +++ b/front/temporal/analytics_queue/signals.ts @@ -7,3 +7,7 @@ import { defineSignal } from "@temporalio/workflow"; export const storeAgentMessageConsumptionAttributionV2Signal = defineSignal< [void] >("store_agent_message_consumption_attribution_v2_signal"); + +export const storeAgentMessageConsumptionAttributionV3Signal = defineSignal< + [void] +>("store_agent_message_consumption_attribution_v3_signal"); diff --git a/front/temporal/analytics_queue/workflows.ts b/front/temporal/analytics_queue/workflows.ts index b757ab75f46c..babb290de8af 100644 --- a/front/temporal/analytics_queue/workflows.ts +++ b/front/temporal/analytics_queue/workflows.ts @@ -1,6 +1,9 @@ import type { AuthenticatorType } from "@app/lib/auth"; import type * as activities from "@app/temporal/analytics_queue/activities"; -import { storeAgentMessageConsumptionAttributionV2Signal } from "@app/temporal/analytics_queue/signals"; +import { + storeAgentMessageConsumptionAttributionV2Signal, + storeAgentMessageConsumptionAttributionV3Signal, +} from "@app/temporal/analytics_queue/signals"; import type { AgentLoopArgs, AgentMessageRef, @@ -21,6 +24,13 @@ const { }, }); +// Consumption indexing is idempotent. The default policy retries without an attempt limit. +const { storeAgentMessageConsumptionAnalyticsActivity } = proxyActivities< + typeof activities +>({ + startToCloseTimeout: "5 minutes", +}); + export async function storeAgentAnalyticsWorkflow( authType: AuthenticatorType, { @@ -86,3 +96,32 @@ export async function storeAgentMessageConsumptionAttributionV2Workflow( }); } } + +// V3 adds consumption analytics indexation after each committed attribution pass. V2 remains +// unchanged above so executions started before this deployment can replay deterministically. +export async function storeAgentMessageConsumptionAttributionV3Workflow( + authType: AuthenticatorType, + { + agentLoopArgs, + }: { + agentLoopArgs: AgentLoopArgs; + } +): Promise { + let pendingRecompute = true; + + setHandler(storeAgentMessageConsumptionAttributionV3Signal, () => { + pendingRecompute = true; + }); + + while (pendingRecompute) { + pendingRecompute = false; + + await storeAgentMessageConsumptionAttributionActivity(authType, { + agentLoopArgs, + }); + + await storeAgentMessageConsumptionAnalyticsActivity(authType, { + agentLoopArgs, + }); + } +}