Skip to content
Open
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
8 changes: 8 additions & 0 deletions .eslintrc.js
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,14 @@ module.exports = {
'node': true,
'jest': true
},
globals: {
/**
* WHATWG Fetch/Streams API globals available in Node 18+ (this project runs on Node 24
* per .nvmrc) - not part of eslint's "node" env, which predates them
*/
'ReadableStream': 'readonly',
'Response': 'readonly'
},
rules: {
'@typescript-eslint/camelcase': 'warn',
'@typescript-eslint/no-unused-vars': 'warn',
Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "hawk.api",
"version": "1.5.11",
"version": "1.5.12",
"main": "index.ts",
"license": "BUSL-1.1",
"scripts": {
Expand Down
2 changes: 1 addition & 1 deletion src/directives/requireUserInWorkspace.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ async function checkUserInWorkspaceByWorkspaceId(context: ResolverContextBase, w
* @param context - request context
* @param projectId - project id
*/
async function checkUserInWorkspaceByProjectId(context: ResolverContextBase, projectId: string): Promise<void> {
export async function checkUserInWorkspaceByProjectId(context: ResolverContextBase, projectId: string): Promise<void> {
const userId = context.user.id;

if (userId) {
Expand Down
6 changes: 6 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import ReleasesFactory from './models/releasesFactory';
import RedisHelper from './redisHelper';
import { appendSsoRoutes } from './sso';
import { appendGitHubRoutes } from './integrations/github';
import { appendAiAssistantRoutes } from './integrations/vercel-ai/routes';

/**
* Option to enable playground
Expand Down Expand Up @@ -272,6 +273,11 @@ class HawkAPI {
*/
appendGitHubRoutes(this.app, sharedFactories);

/**
* Append AI assistant route to Express app
*/
appendAiAssistantRoutes(this.app);

await this.server.start();
this.app.use(graphqlUploadExpress());
this.server.applyMiddleware({ app: this.app });
Expand Down
37 changes: 31 additions & 6 deletions src/integrations/vercel-ai/index.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { generateText } from 'ai';
import { generateText, streamText } from 'ai';
import { ProviderOptions } from '@ai-sdk/provider-utils';

/**
* Params for a single completion call to the model
Expand Down Expand Up @@ -29,11 +30,24 @@ class VercelAIApi {
*/
private readonly modelId: string;

/**
* Provider Gateway fallback order
*/
private readonly providerOptions: ProviderOptions;

/**
* Set up model id and provider fallback order
*/
constructor() {
/**
* @todo make it dynamic, get from project settings
*/
this.modelId = 'deepseek/deepseek-v4-flash';
this.providerOptions = {
gateway: {
order: ['novita', 'azure', 'deepseek'],
},
};
}

/**
Expand All @@ -47,15 +61,26 @@ class VercelAIApi {
model: this.modelId,
system,
prompt,
providerOptions: {
gateway: {
order: ['novita', 'azure', 'deepseek'],
},
},
providerOptions: this.providerOptions,
});

return text;
}

/**
* Send a system/prompt pair to the model and return the generated text as a stream
*
* @param {CompletionParams} params - system instruction and prompt to complete
* @returns {StreamTextResult} text generated by the model, as a stream
*/
public stream({ system, prompt }: CompletionParams): ReturnType<typeof streamText> {
return streamText({
model: this.modelId,
system,
prompt,
providerOptions: this.providerOptions,
});
}
}

export const vercelAIApi = new VercelAIApi();
120 changes: 120 additions & 0 deletions src/integrations/vercel-ai/routes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
import '../../typeDefs/expressContext';
import express from 'express';
import { Readable } from 'stream';
import type { ReadableStream as NodeReadableStream } from 'stream/web';
import { getEventsFactory } from '../../resolvers/helpers/eventsFactory';
import { checkUserInWorkspaceByProjectId } from '../../directives/requireUserInWorkspace';
import { aiService } from '../../services/ai';

/**
* Verify the requesting user is a member of the project's workspace.
*
* @param req - Express request
* @param res - Express response
* @param projectId - project id from query parameters
* @returns user ID if authorized, {@code null} otherwise (response already sent)
*/
async function authorizeProjectAccess(
req: express.Request,
res: express.Response,
projectId: string | undefined
): Promise<string | null> {
const userId = req.context?.user?.id;

if (!userId) {
res.status(401).json({ error: 'Unauthorized. Please provide authorization token.' });

return null;
}

if (!projectId) {
res.status(400).json({ error: 'projectId query parameter is required' });

return null;
}

try {
await checkUserInWorkspaceByProjectId(req.context, projectId);
} catch (error) {
res.status(403).json({ error: error instanceof Error ? error.message : 'You have no access to this workspace' });

return null;
}

return userId;
}

/**
* Create AI assistant router
*
* @returns Express router with AI assistant endpoints
*/
export function createAiStreamRouter(): express.Router {
const router = express.Router();

/**
* GET /integration/ai/stream?projectId=<projectId>&eventId=<eventId>&originalEventId=<originalEventId>
* Stream an AI suggestion for the event
*/
router.get('/stream', async (req, res, next) => {
try {
const { projectId, eventId, originalEventId } = req.query;

const userId = await authorizeProjectAccess(req, res, projectId as string | undefined);

if (!userId) {
return;
}

if (!eventId || typeof eventId !== 'string') {
res.status(400).json({ error: 'eventId query parameter is required' });

return;
}

if (!originalEventId || typeof originalEventId !== 'string') {
res.status(400).json({ error: 'originalEventId query parameter is required' });

return;
}

const eventsFactory = getEventsFactory(req.context, projectId as string);

let result;

try {
result = await aiService.streamSuggestion(eventsFactory, eventId, originalEventId);
} catch (error) {
res.status(404).json({ error: error instanceof Error ? error.message : 'Event not found' });

return;
}

const response = result.toUIMessageStreamResponse();

res.status(response.status);
response.headers.forEach((value, key) => res.setHeader(key, value));

if (!response.body) {
res.end();

return;
}

Readable.fromWeb(response.body as NodeReadableStream<Uint8Array>).pipe(res);
} catch (error) {
next(error);
}
});
Comment on lines +93 to +108

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we can just use

result.pipeUIMessageStreamToResponse(res);

with extra options or just

result.pipeTextStreamToResponse(res);

since we don't need any metadata toolcalls etc


return router;
}

/**
* Append AI assistant routes to Express app
*
* @param app - Express application instance
*/
export function appendAiAssistantRoutes(app: express.Application): void {
app.use('/integration/ai', createAiStreamRouter());
}
61 changes: 55 additions & 6 deletions src/services/ai.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { buildEventPrompt, spotlightInstruction } from './askAi/security/spotlig
import { isLeaked, SUGGESTION_FALLBACK_MESSAGE } from './askAi/security/leakDetector';
import { ctoInstruction } from './askAi/instructions/cto';
import { EventsFactoryInterface } from './types';
import type { Event } from './types';

/**
* Report that the leak tripwire fired.
Expand Down Expand Up @@ -43,12 +44,12 @@ export class AIService {
* @param originalEventId - original event id
* @returns {Promise<string>} - suggestion
*/
public async generateSuggestion(eventsFactory: EventsFactoryInterface, eventId: string, originalEventId: string): Promise<string> {
const event = await eventsFactory.getEventRepetition(eventId, originalEventId);

if (!event) {
throw new Error('Event not found');
}
public async generateSuggestion(
eventsFactory: EventsFactoryInterface,
eventId: string,
originalEventId: string
): Promise<string> {
const event = await this.getEventOrThrow(eventsFactory, eventId, originalEventId);

const { prompt, nonce } = buildEventPrompt(event.payload);

Expand All @@ -65,6 +66,54 @@ export class AIService {

return text;
}

/**
* Generate streaming suggestion for the event
*
* The payload is spotlighted by {@link buildEventPrompt} exactly as in
* {@link AIService.generateSuggestion}.
*
* @param eventsFactory - events factory
* @param eventId - event id
* @param originalEventId - original event id
* @returns streaming suggestion
*/
public async streamSuggestion(
eventsFactory: EventsFactoryInterface,
eventId: string,
originalEventId: string
): Promise<ReturnType<typeof vercelAIApi.stream>> {
const event = await this.getEventOrThrow(eventsFactory, eventId, originalEventId);

const { prompt, nonce } = buildEventPrompt(event.payload);

return vercelAIApi.stream({
system: ctoInstruction + spotlightInstruction(nonce),
prompt,
});
}

/**
* Find the event repetition or throw if it doesn't exist
*
* @param eventsFactory - events factory
* @param eventId - event id
* @param originalEventId - original event id
* @returns {Promise<Event>} - event repetition
*/
private async getEventOrThrow(
eventsFactory: EventsFactoryInterface,
eventId: string,
originalEventId: string
): Promise<Event> {
const event = await eventsFactory.getEventRepetition(eventId, originalEventId);

if (!event) {
throw new Error('Event not found');
}

return event;
}
}

export const aiService = new AIService();
4 changes: 2 additions & 2 deletions src/services/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ import { EventAddons, EventData } from '@hawk.so/types';
/**
* Event type which is returned by events factory
*/
type Event = {
export type Event = {
_id: string;
payload: EventData<EventAddons>;
};
Expand All @@ -20,4 +20,4 @@ export interface EventsFactoryInterface {
* @returns {Promise<EventData<EventAddons>>} - event repetition
*/
getEventRepetition(repetitionId: string, originalEventId: string): Promise<Event>;
}
}
Loading
Loading