fix(sync): close local queue semantic merge gap
ci/woodpecker/pr/ci Pipeline was successful

This commit is contained in:
2026-08-02 23:16:09 -05:00
parent 906ad8dc30
commit b8844e1ff0
5 changed files with 62 additions and 2 deletions
@@ -72,13 +72,13 @@ const mockChatGateway = {
broadcastSessionInfo: vi.fn(), broadcastSessionInfo: vi.fn(),
}; };
function buildService(): CommandExecutorService { function buildService(redis: typeof mockRedis | null = mockRedis): CommandExecutorService {
return new CommandExecutorService( return new CommandExecutorService(
mockRegistry as never, mockRegistry as never,
mockAgentService as never, mockAgentService as never,
mockSystemOverride as never, mockSystemOverride as never,
mockSessionGC as never, mockSessionGC as never,
mockRedis as never, redis as never,
mockBrain as never, mockBrain as never,
null, null,
mockChatGateway as never, mockChatGateway as never,
@@ -131,6 +131,22 @@ describe('CommandExecutorService — P8-012 commands', () => {
expect(ttl).toBe(300); expect(ttl).toBe(300);
}); });
it('/provider login remains available without Redis on the local tier', async () => {
const localService = buildService(null);
const payload: SlashCommandPayload = {
command: 'provider',
args: 'login anthropic',
conversationId,
};
const result = await localService.execute(payload, userScope);
expect(result.success).toBe(true);
expect(result.message).not.toContain('token=');
expect(result.data).toEqual({ provider: 'anthropic' });
expect(mockRedis.set).not.toHaveBeenCalled();
});
// /provider with no args — returns usage // /provider with no args — returns usage
it('/provider with no args returns usage message', async () => { it('/provider with no args returns usage message', async () => {
const payload: SlashCommandPayload = { command: 'provider', conversationId }; const payload: SlashCommandPayload = { command: 'provider', conversationId };
@@ -119,6 +119,19 @@ describe('SessionGCService', () => {
).resolves.toEqual({ allowed: true }); ).resolves.toEqual({ allowed: true });
}); });
it('collect() skips Valkey but still demotes only the requested session on local tier', async () => {
const localService = new SessionGCService(null, mockLogService as unknown as LogService);
const result = await localService.collect('local-session');
expect(result.sessionId).toBe('local-session');
expect(result.cleaned.valkeyKeys).toBeUndefined();
expect(mockLogService.logs.promoteSessionToWarm).toHaveBeenCalledWith(
'local-session',
expect.any(Date),
);
});
it('collect() returns sessionId in result', async () => { it('collect() returns sessionId in result', async () => {
const result = await service.collect('test-session-id'); const result = await service.collect('test-session-id');
expect(result.sessionId).toBe('test-session-id'); expect(result.sessionId).toBe('test-session-id');
@@ -0,0 +1,23 @@
import { describe, expect, it } from 'vitest';
import type { MosaicConfig } from '@mosaicstack/config';
import { SystemOverrideService } from './system-override.service.js';
const localConfig = { queue: { type: 'local' } } as MosaicConfig;
describe('SystemOverrideService local tier', () => {
it('keeps ephemeral overrides isolated by tenant and user scope', async () => {
const service = new SystemOverrideService(localConfig);
const firstScope = { tenantId: 'tenant-a', userId: 'user-a' };
const secondScope = { tenantId: 'tenant-b', userId: 'user-b' };
await service.set('shared-session', 'first override', firstScope);
await service.set('shared-session', 'second override', secondScope);
await expect(service.get('shared-session', firstScope)).resolves.toBe('first override');
await expect(service.get('shared-session', secondScope)).resolves.toBe('second override');
await service.clear('shared-session', firstScope);
await expect(service.get('shared-session', firstScope)).resolves.toBeNull();
await expect(service.get('shared-session', secondScope)).resolves.toBe('second override');
});
});
@@ -17,6 +17,7 @@ describe('QueueService local tier', () => {
await expect( await expect(
service.addRepeatableJob('mosaic-test', 'local-noop', {}, '* * * * *'), service.addRepeatableJob('mosaic-test', 'local-noop', {}, '* * * * *'),
).resolves.toBeUndefined(); ).resolves.toBeUndefined();
await expect(service.removeRepeatableJobs('mosaic-test', 'local-noop')).resolves.toBe(0);
await expect(service.getHealthStatus()).resolves.toEqual({ queues: {}, healthy: true }); await expect(service.getHealthStatus()).resolves.toEqual({ queues: {}, healthy: true });
await expect(service.listJobs()).resolves.toEqual([]); await expect(service.listJobs()).resolves.toEqual([]);
await expect(service.retryJob('mosaic-test__1')).resolves.toEqual({ await expect(service.retryJob('mosaic-test__1')).resolves.toEqual({
+7
View File
@@ -199,7 +199,14 @@ export class QueueService implements OnModuleInit, OnModuleDestroy {
* safe retirement of previously registered system-wide jobs. * safe retirement of previously registered system-wide jobs.
*/ */
async removeRepeatableJobs(queueName: string, jobName: string): Promise<number> { async removeRepeatableJobs(queueName: string, jobName: string): Promise<number> {
if (!this.enabled) {
this.logger.debug(
`Skipping repeatable-job removal for "${jobName}" on "${queueName}" (local tier — BullMQ disabled)`,
);
return 0;
}
const queue = this.getQueue(queueName); const queue = this.getQueue(queueName);
if (!queue) return 0;
const jobs = await queue.getRepeatableJobs(); const jobs = await queue.getRepeatableJobs();
const matchingJobs = jobs.filter((job) => job.name === jobName); const matchingJobs = jobs.filter((job) => job.name === jobName);
await Promise.all(matchingJobs.map((job) => queue.removeRepeatableByKey(job.key))); await Promise.all(matchingJobs.map((job) => queue.removeRepeatableByKey(job.key)));