Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
105 changes: 105 additions & 0 deletions packages/apps/src/app.process.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { context, propagation, ROOT_CONTEXT } from '@opentelemetry/api';
import type { Baggage, Context, ContextManager, Span, Tracer } from '@opentelemetry/api';

import { IMessageActivity, InvokeResponse, ISignInFailureInvokeActivity, ITaskFetchInvokeActivity, IToken, MessageActivity, TaskModuleResponse } from '@microsoft/teams.api';
import { IStorage } from '@microsoft/teams.common';

import { ActivitySender } from './activity-sender';
import { App } from './app';
Expand All @@ -18,6 +19,7 @@ import {
} from './diagnostics/helpers';
import { IActivityResponseEvent, IActivitySentEvent, IErrorEvent } from './events';
import { IActivityEvent } from './events/activity';
import { TurnStateContainer } from './state';
import { createTestApp } from './test-utils';

jest.mock('./diagnostics/helpers', () => ({
Expand Down Expand Up @@ -135,6 +137,109 @@ describe('App', () => {
});

describe('process', () => {
it('loads, exposes, persists, and seals per-turn state', async () => {
const data = new Map<string, string>();
const storage: IStorage<string, string> = {
get: (key) => data.get(key),
set: (key, value) => {
data.set(key, value);
},
delete: (key) => {
data.delete(key);
},
};
await app.stop();
app = createTestApp({ state: { storage } });
await app.start();
const turnStates: TurnStateContainer[] = [];
const counts: number[] = [];
app.on('message', ({ state }) => {
if (!state) {
throw new Error('Expected state to be enabled.');
}
turnStates.push(state);
const count = (state.conversation.get<number>('count') ?? 0) + 1;
counts.push(count);
state.conversation.set('count', count);
});
const stateActivity = new MessageActivity('hello')
.withFrom({ id: 'user-1', name: 'Test User', role: 'user' })
.withRecipient({ id: 'bot-1', name: 'Test Bot', role: 'bot' })
.withConversation({ id: 'conv-1', conversationType: 'personal' })
.withChannelId('msteams')
.toInterface();

await app.process({ token, body: stateActivity });
await app.process({ token, body: stateActivity });

expect(turnStates).toHaveLength(2);
expect(counts).toEqual([1, 2]);
expect(turnStates[0].conversation.isSealed).toBe(true);
expect(() => turnStates[0].conversation.get('count')).toThrow();
});

it('persists dirty state when a handler fails', async () => {
const data = new Map<string, string>();
const storage: IStorage<string, string> = {
get: (key) => data.get(key),
set: (key, value) => {
data.set(key, value);
},
delete: (key) => {
data.delete(key);
},
};
await app.stop();
app = createTestApp({ state: { storage } });
await app.start();
app.on('message', ({ state }) => {
state?.conversation.set('saved', true);
throw new Error('handler failed');
});
const stateActivity = new MessageActivity('hello')
.withFrom({ id: 'user-1', name: 'Test User', role: 'user' })
.withRecipient({ id: 'bot-1', name: 'Test Bot', role: 'bot' })
.withConversation({ id: 'conv-1', conversationType: 'personal' })
.withChannelId('msteams')
.toInterface();

const response = await app.process({ token, body: stateActivity });

expect(response.status).toBe(500);
expect(JSON.parse(data.get('ts:conv:conv-1') ?? '{}').data).toEqual({
saved: true,
});
});

it('seals state and propagates the error when persistence fails', async () => {
const saveError = new Error('save failed');
const storage: IStorage<string, string> = {
get: () => undefined,
set: () => {
throw saveError;
},
delete: () => undefined,
};
await app.stop();
app = createTestApp({ state: { storage } });
await app.start();
let capturedState: TurnStateContainer | undefined;
app.on('message', ({ state }) => {
capturedState = state;
state?.conversation.set('saved', true);
});
const stateActivity = new MessageActivity('hello')
.withFrom({ id: 'user-1', name: 'Test User', role: 'user' })
.withRecipient({ id: 'bot-1', name: 'Test Bot', role: 'bot' })
.withConversation({ id: 'conv-1', conversationType: 'personal' })
.withChannelId('msteams')
.toInterface();

await expect(app.process({ token, body: stateActivity })).rejects.toBe(saveError);
expect(capturedState?.conversation.isSealed).toBe(true);
expect(() => capturedState?.conversation.get('saved')).toThrow();
});

it('should return status 200 if no route matches', async () => {
const event: IActivityEvent = {
token: token,
Expand Down
21 changes: 21 additions & 0 deletions packages/apps/src/app.process.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import { IActivityEvent } from './events';
import { Router } from './router';
import type { Route } from './router/route';
import { IRoutes } from './routes';
import { TurnStateLoader } from './state';
import { IActivitySender, IPlugin, RouteHandler, StreamCancelledError } from './types';
import { PluginAdditionalContext } from './types/app-routing';

Expand Down Expand Up @@ -75,6 +76,10 @@ export interface IActivityProcessorOptions<TPlugin extends IPlugin = IPlugin> {
readonly api: ApiClient;
readonly client: HttpClient;
readonly storage: IStorage;
/**
* Loader used to attach and persist state for each activity turn.
*/
readonly stateLoader?: TurnStateLoader;
readonly log: ILogger;
readonly getId: () => string | undefined;
readonly getConnectionName: () => string;
Expand Down Expand Up @@ -253,6 +258,14 @@ export class ActivityProcessor<TPlugin extends IPlugin = IPlugin> {
...pluginContexts
});

const conversationId = activity.conversation?.id;
if (this.options.stateLoader && conversationId) {
context.state = await this.options.stateLoader.load(
conversationId,
activity.from?.id
);
}

const send = context.send.bind(context);
context.send = async (activity: ActivityLike | DeprecatedInputActivity, conversationRef?: ConversationReference) => {
const res = await send(activity, conversationRef ?? ref);
Expand Down Expand Up @@ -312,6 +325,14 @@ export class ActivityProcessor<TPlugin extends IPlugin = IPlugin> {
activity,
response: response,
});
} finally {
if (context.state && this.options.stateLoader) {
try {
await this.options.stateLoader.save(context.state);
} finally {
context.state.seal();
}
}
}

return response;
Expand Down
13 changes: 13 additions & 0 deletions packages/apps/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ import { DEFAULT_OAUTH_SETTINGS, OAuthSettings } from './oauth';
import { HttpPlugin } from './plugins';
import { Router } from './router';
import { IRoutes } from './routes';
import { createStateLoader, StateOptions, TurnStateLoader } from './state';
import { DEFAULT_TENANT_FOR_GRAPH_TOKEN, TokenManager } from './token-manager';
import { AppTokenProvider, IAppTokenProvider } from './token-provider';
import { AppEvents, IPlugin, PluginName, RouteHandler } from './types';
Expand Down Expand Up @@ -170,6 +171,15 @@ export type AppOptions<TPlugin extends IPlugin> = {
*/
readonly storage?: IStorage;

/**
* Enables per-turn conversation and user state.
*
* Pass `true` to use the app storage with default keys, or provide options
* to configure dedicated storage, key prefix, or expiration. State is
* disabled when omitted or `false`.
*/
readonly state?: boolean | StateOptions;

/**
* plugins to extend the apps functionality
*/
Expand Down Expand Up @@ -329,6 +339,7 @@ export class App<TPlugin extends IPlugin = IPlugin> {
private readonly tokenManager: TokenManager;

private readonly _tokenProvider: AppTokenProvider;
private readonly stateLoader?: TurnStateLoader;

private eventManager!: EventManager<TPlugin>;
private activityProcessor!: ActivityProcessor<TPlugin>;
Expand All @@ -339,6 +350,7 @@ export class App<TPlugin extends IPlugin = IPlugin> {
constructor(readonly options: AppOptions<TPlugin> = {}) {
this.log = this.options.logger || new ConsoleLogger('@teams/app');
this.storage = this.options.storage || new LocalStorage();
this.stateLoader = createStateLoader(this.options.state, this.storage, this.log);

// Resolve cloud environment from options or CLOUD env var
const cloudEnvName = typeof process !== 'undefined' ? process.env.CLOUD : undefined;
Expand Down Expand Up @@ -446,6 +458,7 @@ export class App<TPlugin extends IPlugin = IPlugin> {
api: this.api,
client: this.client,
storage: this.storage,
stateLoader: this.stateLoader,
log: this.log,
getId: () => this.id,
getConnectionName: () => this.oauth.defaultConnectionName,
Expand Down
11 changes: 11 additions & 0 deletions packages/apps/src/contexts/activity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
import { ILogger, IStorage } from '@microsoft/teams.common';

import { ApiClient, GraphClient } from '../api';
import { TurnStateContainer } from '../state';
import { IStreamer } from '../types';
import { IActivitySender } from '../types/plugin/sender';

Expand Down Expand Up @@ -87,6 +88,14 @@ export interface IBaseActivityContextOptions<T extends Activity = Activity> {
*/
storage: IStorage;

/**
* Conversation and user state loaded for this activity turn.
*
* This is undefined when state is disabled or the activity has no conversation ID.
* Do not retain the container after the handler completes; its scopes are sealed.
*/
state?: TurnStateContainer;

/**
* whether the user has provided
* their MSGraph credentials for use
Expand Down Expand Up @@ -220,6 +229,7 @@ export class ActivityContext<T extends Activity = Activity, TExtraCtx extends {}
appGraph!: GraphClient;
userGraph!: GraphClient;
storage!: IStorage;
state?: TurnStateContainer;
stream!: IStreamer;
isSignedIn?: boolean;
connectionName: string;
Expand Down Expand Up @@ -458,6 +468,7 @@ export class ActivityContext<T extends Activity = Activity, TExtraCtx extends {}
log: this.log,
ref: this.ref,
storage: this.storage,
state: this.state,
stream: this.stream,
isSignedIn: this.isSignedIn,
connectionName: this.connectionName,
Expand Down
1 change: 1 addition & 0 deletions packages/apps/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ export * from './types';
export * from './contexts';
export * from './oauth';
export * from './events';
export * from './state';
// Only the interface is public: the implementing class is constructed from the
// app's `TokenManager`, which is internal.
export type { IAppTokenProvider } from './token-provider';
Expand Down
63 changes: 63 additions & 0 deletions packages/apps/src/state/container.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import { TurnState } from './turn-state';

type StateDeleter = (conversationId: string, userId?: string) => Promise<void>;

/** Conversation and user state loaded for one activity turn. */
export class TurnStateContainer {
/** State shared by all users in the current conversation. */
readonly conversation: TurnState;

/**
* State for the current user within the current conversation.
* Undefined when the activity has no sender ID.
*/
readonly user?: TurnState;

/** Conversation ID used to load and persist this container. */
readonly conversationId: string;

/** User ID used to load and persist the user scope. */
readonly userId?: string;

private readonly deleter: StateDeleter;

/**
* Creates a loaded state container.
* @param conversation Conversation-scoped state.
* @param conversationId Conversation ID associated with the state.
* @param deleter Callback that removes persisted scopes.
* @param user Optional user-within-conversation state.
* @param userId Optional sender ID associated with the user scope.
*/
constructor(
conversation: TurnState,
conversationId: string,
deleter: StateDeleter,
user?: TurnState,
userId?: string
) {
this.conversation = conversation;
this.conversationId = conversationId;
this.deleter = deleter;
this.user = user;
this.userId = userId;
}

/**
* Deletes both persisted scopes and clears their in-memory snapshots.
*
* The backing delete must succeed before the in-memory state changes. Values
* written after this call are persisted normally at the end of the turn.
*/
async delete(): Promise<void> {
await this.deleter(this.conversationId, this.userId);
this.conversation.reset();
this.user?.reset();
}

/** Seals both scopes after activity processing completes. */
seal(): void {
this.conversation.seal();
this.user?.seal();
}
}
4 changes: 4 additions & 0 deletions packages/apps/src/state/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
export { TurnState, TurnStateSealedError } from './turn-state';
export { TurnStateContainer } from './container';
export { TurnStateLoader, createStateLoader } from './loader';
export type { StateOptions } from './options';
Loading
Loading