Skip to content
Merged
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
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
import { beforeEach, describe, expect, it } from 'vitest';

import { activeProgressTargetStore, getActiveProgressTargets } from './activeProgressTargetStore';
import {
activeProgressTargetStore,
getActiveProgressTargets,
getFollowedProgressTargets,
} from './activeProgressTargetStore';

const target = (queueItemId: string, itemIndex: number) => ({ itemIndex, queueItemId });

Expand Down Expand Up @@ -65,4 +69,50 @@ describe('activeProgressTargetStore', () => {

expect(getActiveProgressTargets()[0]).toEqual(target('queue-1', 1));
});

it('keeps a settling slot followable but out of the running set', () => {
// Completed, result not routed yet: the single-slot preview keeps following
// it, while the tile grid must not count it or a single-GPU batch would
// flash into two tiles at every item boundary.
activeProgressTargetStore.set(target('queue-1', 1));
activeProgressTargetStore.set(target('queue-1', 2));

activeProgressTargetStore.settle(target('queue-1', 1));

expect(getActiveProgressTargets()).toEqual([target('queue-1', 2)]);
// Running first: the settling slot is followed only when nothing is running.
expect(getFollowedProgressTargets()).toEqual([target('queue-1', 2), target('queue-1', 1)]);

activeProgressTargetStore.clear(target('queue-1', 1));

expect(getFollowedProgressTargets()).toEqual([target('queue-1', 2)]);
});

it('ignores settling a slot that never reported progress', () => {
activeProgressTargetStore.set(target('queue-1', 1));
const first = getFollowedProgressTargets();

activeProgressTargetStore.settle(target('queue-1', 2));

expect(getFollowedProgressTargets()).toBe(first);
});

it('returns a settled slot to the running set when it reports progress again', () => {
activeProgressTargetStore.set(target('queue-1', 1));
activeProgressTargetStore.settle(target('queue-1', 1));

activeProgressTargetStore.set(target('queue-1', 1));

expect(getActiveProgressTargets()).toEqual([target('queue-1', 1)]);
expect(getFollowedProgressTargets()).toEqual([target('queue-1', 1)]);
});

it('clears settling slots along with running ones', () => {
activeProgressTargetStore.set(target('queue-1', 1));
activeProgressTargetStore.settle(target('queue-1', 1));

activeProgressTargetStore.clear();

expect(getFollowedProgressTargets()).toEqual([]);
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { registerAccountOwnedResource } from '@platform/state/accountLifecycle';
import { createExternalStore } from '@platform/state/externalStore';

/**
* The slots currently reporting progress.
* The slots currently reporting progress, plus the ones settling.
*
* A list rather than a single value because of multi-GPU: with `generation_devices`
* (default `auto`) the backend runs one session per GPU, so a batch of four across
Expand All @@ -14,46 +14,83 @@ import { createExternalStore } from '@platform/state/externalStore';
*
* Order is the order sessions started, which keeps the single-target accessor below
* stable for as long as that session runs.
*
* A *settling* slot is one whose backend item has completed but whose result has
* not landed in the gallery yet — two HTTP round trips away. Single-slot surfaces
* keep following it so the last denoise frame stays up until the finished image
* can take over; multi-slot surfaces (the tile grid) stop counting it, or a
* single-GPU batch would flash into a two-tile grid at every item boundary.
*/
export interface ActiveProgressTargetSink {
clear(target?: QueueItemProgressTarget): void;
set(target: QueueItemProgressTarget): void;
settle(target: QueueItemProgressTarget): void;
}

interface ActiveProgressTargetsSnapshot {
settlingTargets: QueueItemProgressTarget[];
targets: QueueItemProgressTarget[];
}

const store = createExternalStore<{ targets: QueueItemProgressTarget[] }>({ targets: [] });
const store = createExternalStore<ActiveProgressTargetsSnapshot>({ settlingTargets: [], targets: [] });

const isSameTarget = (left: QueueItemProgressTarget, right: QueueItemProgressTarget): boolean =>
left.queueItemId === right.queueItemId && left.itemIndex === right.itemIndex;

const includes = (targets: QueueItemProgressTarget[], target: QueueItemProgressTarget): boolean =>
targets.some((candidate) => isSameTarget(candidate, target));

const without = (targets: QueueItemProgressTarget[], target: QueueItemProgressTarget): QueueItemProgressTarget[] =>
targets.filter((candidate) => !isSameTarget(candidate, target));

export const activeProgressTargetStore: ActiveProgressTargetSink = {
clear(target) {
const { targets } = store.getSnapshot();
const { settlingTargets, targets } = store.getSnapshot();

if (!target) {
if (targets.length > 0) {
store.patchSnapshot({ targets: [] });
if (targets.length > 0 || settlingTargets.length > 0) {
store.patchSnapshot({ settlingTargets: [], targets: [] });
}

return;
}

const remaining = targets.filter((candidate) => !isSameTarget(candidate, target));
const remaining = without(targets, target);
const remainingSettling = without(settlingTargets, target);

if (remaining.length !== targets.length) {
store.patchSnapshot({ targets: remaining });
}
store.patchSnapshot({
...(remaining.length !== targets.length ? { targets: remaining } : {}),
...(remainingSettling.length !== settlingTargets.length ? { settlingTargets: remainingSettling } : {}),
});
},
set(target) {
const { targets } = store.getSnapshot();
const { settlingTargets, targets } = store.getSnapshot();

// Progress frames arrive many times a second per session; re-appending an
// already-tracked target would publish a fresh array identity every frame and
// re-render every consumer.
if (targets.some((candidate) => isSameTarget(candidate, target))) {
if (includes(targets, target)) {
return;
}

store.patchSnapshot({
targets: [...targets, target],
// A settled slot reporting progress again is running again.
...(includes(settlingTargets, target) ? { settlingTargets: without(settlingTargets, target) } : {}),
});
},
settle(target) {
const { settlingTargets, targets } = store.getSnapshot();

// A slot that never reported progress was never followed; nothing to keep up.
if (!includes(targets, target)) {
return;
}

store.patchSnapshot({ targets: [...targets, target] });
store.patchSnapshot({
settlingTargets: includes(settlingTargets, target) ? settlingTargets : [...settlingTargets, target],
targets: without(targets, target),
});
},
};

Expand All @@ -62,19 +99,36 @@ registerAccountOwnedResource({
name: 'queue-active-progress-target',
});

/**
* Running slots first: a settling slot is only worth following while nothing is
* running, or a concurrent session's live stream would sit unseen behind a
* static frame for the whole routing window.
*/
const selectFollowedTargets = ({
settlingTargets,
targets,
}: ActiveProgressTargetsSnapshot): QueueItemProgressTarget[] =>
settlingTargets.length === 0 ? targets : [...targets, ...settlingTargets];

/**
* The slot to follow where a surface can only show one.
*
* The oldest still-running slot rather than the most recent to report: following the
* most recent is what made the preview flip between concurrent sessions. Behaviour is
* identical to the previous single-value store whenever one session runs at a time,
* which is every single-GPU install.
* which is every single-GPU install — except that a completed slot stays followed
* until its result lands.
*/
export const useActiveProgressTarget = (): QueueItemProgressTarget | null =>
store.useSelector((snapshot) => snapshot.targets[0] ?? null);
store.useSelector((snapshot) => selectFollowedTargets(snapshot)[0] ?? null);

/** Every slot reporting progress, in the order its session started. */
/** Every slot currently running, in the order its session started. */
export const useActiveProgressTargets = (): QueueItemProgressTarget[] =>
store.useSelector((snapshot) => snapshot.targets);

/** Every followable slot — running ones first, then settling ones. */
export const useFollowedProgressTargets = (): QueueItemProgressTarget[] => store.useSelector(selectFollowedTargets);

export const getActiveProgressTargets = (): QueueItemProgressTarget[] => store.getSnapshot().targets;

export const getFollowedProgressTargets = (): QueueItemProgressTarget[] => selectFollowedTargets(store.getSnapshot());
6 changes: 6 additions & 0 deletions invokeai/frontend/webv2/src/features/queue/data/events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,12 @@ export interface InvocationProgressEvent extends InvocationEventBase {
percentage: number | null;
/** Intermittent denoising preview, when the invocation produces one. */
image?: { width: number; height: number; dataURL: string } | null;
/**
* Monotonic per queue item, when the backend sends it: a frame at or below a
* revision already shown is stale and dropped. Absent from today's socket
* events, which arrive in order; the reconnect snapshot carries it.
*/
revision?: number | null;
/**
* The accelerator running this session, e.g. `cuda:1` or `xpu:1` — null on CPU/MPS and in
* single-device mode. With `generation_devices` set (default `auto`) several
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';

import {
consumeQueueItemSwapProgressImage,
getQueueItemBridgeProgressImage,
getQueueItemSwapProgressImage,
progressImageStore,
SWAP_FRAME_TTL_MS,
} from './progressImageStore';

const frame = (label: string) => ({ dataUrl: `data:image/png;base64,${label}`, height: 32, width: 64 });
const target = (queueItemId: string, itemIndex = 1) => ({ itemIndex, queueItemId });

describe('progressImageStore held frames', () => {
beforeEach(() => {
vi.useFakeTimers();
progressImageStore.clear();
});

afterEach(() => {
progressImageStore.clear();
vi.useRealTimers();
});

it('copies the slot frame into both held sets and keeps the live frame', () => {
progressImageStore.set(frame('last'), target('queue-1'));

progressImageStore.hold(target('queue-1'));
progressImageStore.bindSwapImages('queue-1', ['result.png']);

expect(getQueueItemBridgeProgressImage('queue-1')).toEqual(frame('last'));
expect(getQueueItemSwapProgressImage('queue-1', 'result.png')).toEqual(frame('last'));
// The single-slot preview keeps showing the live frame while routing runs.
progressImageStore.clear(target('queue-1'));
expect(getQueueItemBridgeProgressImage('queue-1')).toEqual(frame('last'));
});

it('holds nothing for a slot that never produced a frame', () => {
progressImageStore.hold(target('queue-1'));
progressImageStore.bindSwapImages('queue-1', ['result.png']);

expect(getQueueItemBridgeProgressImage('queue-1')).toBeNull();
expect(getQueueItemSwapProgressImage('queue-1', 'result.png')).toBeNull();
});

it('paints the swap frame only over the images its backend item delivered', () => {
// Item 3 of a batch finishing must not put its denoise frame over item 1's
// image when the user clicks that one, nor over anything before routing
// has said which images the frame belongs to.
progressImageStore.set(frame('third'), target('queue-1', 3));
progressImageStore.hold(target('queue-1', 3));

expect(getQueueItemSwapProgressImage('queue-1', 'image-3.png')).toBeNull();

progressImageStore.bindSwapImages('queue-1', ['image-3.png', 'image-3-control.png']);

expect(getQueueItemSwapProgressImage('queue-1', 'image-3.png')).toEqual(frame('third'));
expect(getQueueItemSwapProgressImage('queue-1', 'image-3-control.png')).toEqual(frame('third'));
expect(getQueueItemSwapProgressImage('queue-1', 'image-1.png')).toBeNull();

// A later hold for the same queue item starts unbound again.
progressImageStore.set(frame('fourth'), target('queue-1', 4));
progressImageStore.hold(target('queue-1', 4));

expect(getQueueItemSwapProgressImage('queue-1', 'image-3.png')).toBeNull();
});

it('consumes the swap frame alone once the finished image has decoded', () => {
progressImageStore.set(frame('last'), target('queue-1'));
progressImageStore.hold(target('queue-1'));
progressImageStore.bindSwapImages('queue-1', ['result.png']);

consumeQueueItemSwapProgressImage('queue-1');

expect(getQueueItemSwapProgressImage('queue-1', 'result.png')).toBeNull();
// The bridge to the batch's next slot is still needed.
expect(getQueueItemBridgeProgressImage('queue-1')).toEqual(frame('last'));
});

it('expires the swap frame so browsing back never replays the low-resolution frame', () => {
progressImageStore.set(frame('last'), target('queue-1'));
progressImageStore.hold(target('queue-1'));
progressImageStore.bindSwapImages('queue-1', ['result.png']);

vi.advanceTimersByTime(SWAP_FRAME_TTL_MS - 1);
expect(getQueueItemSwapProgressImage('queue-1', 'result.png')).toEqual(frame('last'));

vi.advanceTimersByTime(1);
expect(getQueueItemSwapProgressImage('queue-1', 'result.png')).toBeNull();
expect(getQueueItemBridgeProgressImage('queue-1')).toEqual(frame('last'));
});

it('restarts the expiry when a later slot of the same queue item is held', () => {
progressImageStore.set(frame('first'), target('queue-1', 1));
progressImageStore.hold(target('queue-1', 1));
vi.advanceTimersByTime(SWAP_FRAME_TTL_MS - 1);

progressImageStore.set(frame('second'), target('queue-1', 2));
progressImageStore.hold(target('queue-1', 2));
progressImageStore.bindSwapImages('queue-1', ['second.png']);
vi.advanceTimersByTime(SWAP_FRAME_TTL_MS - 1);

expect(getQueueItemSwapProgressImage('queue-1', 'second.png')).toEqual(frame('second'));
});

it('keeps only the most recent queue items', () => {
for (let index = 0; index < 9; index += 1) {
progressImageStore.set(frame(`frame-${index}`), target(`queue-${index}`));
progressImageStore.hold(target(`queue-${index}`));
progressImageStore.bindSwapImages(`queue-${index}`, [`image-${index}.png`]);
}

expect(getQueueItemBridgeProgressImage('queue-0')).toBeNull();
expect(getQueueItemSwapProgressImage('queue-0', 'image-0.png')).toBeNull();
expect(getQueueItemBridgeProgressImage('queue-1')).toEqual(frame('frame-1'));
expect(getQueueItemSwapProgressImage('queue-8', 'image-8.png')).toEqual(frame('frame-8'));
});

it('forgets a queue item on clearHeld and everything on clear', () => {
progressImageStore.set(frame('one'), target('queue-1'));
progressImageStore.hold(target('queue-1'));
progressImageStore.set(frame('two'), target('queue-2'));
progressImageStore.hold(target('queue-2'));
progressImageStore.bindSwapImages('queue-2', ['two.png']);

progressImageStore.clearHeld('queue-1');
expect(getQueueItemBridgeProgressImage('queue-1')).toBeNull();
expect(getQueueItemSwapProgressImage('queue-2', 'two.png')).toEqual(frame('two'));

progressImageStore.clear();
expect(getQueueItemBridgeProgressImage('queue-2')).toBeNull();
expect(getQueueItemSwapProgressImage('queue-2', 'two.png')).toBeNull();
expect(vi.getTimerCount()).toBe(0);
});
});
Loading
Loading