Convert to batch queries.

This commit is contained in:
Dermot Duffy
2023-01-24 19:36:54 -08:00
parent 92b75c4b1a
commit 04d3ca3f46
10 changed files with 600 additions and 345 deletions
+44 -23
View File
@@ -1,35 +1,32 @@
import { CameraConfig } from '../types';
import { RecordingSegmentsCache } from './cache';
import { RecordingSegmentsCache, RequestCache } from './cache';
import { CameraManagerEngine } from './engine';
import { FrigateCameraManagerEngine } from './frigate/engine-frigate';
import { DataQuery } from './types';
import { Engine } from './types';
type CameraManagerEngineCameraIDMap = Map<CameraManagerEngine, Set<string>>;
export class CameraManagerEngineFactory {
protected _engines: Map<string, CameraManagerEngine> = new Map();
protected _engines: Map<Engine, CameraManagerEngine> = new Map();
protected _getOrCreateEngine(engineKey: string): CameraManagerEngine | null {
const cachedEngine = this._engines.get(engineKey);
public getEngine(engine: Engine): CameraManagerEngine | null {
const cachedEngine = this._engines.get(engine);
if (cachedEngine) {
return cachedEngine;
}
let newEngine: CameraManagerEngine | null = null;
switch (engineKey) {
case 'frigate':
newEngine = new FrigateCameraManagerEngine(new RecordingSegmentsCache());
let cameraManagerEngine: CameraManagerEngine | null = null;
switch (engine) {
case Engine.Frigate:
cameraManagerEngine = new FrigateCameraManagerEngine(
new RecordingSegmentsCache(),
new RequestCache(),
);
break;
}
if (newEngine) {
this._engines.set(engineKey, newEngine);
if (cameraManagerEngine) {
this._engines.set(engine, cameraManagerEngine);
}
return newEngine;
}
public getEngineForQuery(
cameras: Map<string, CameraConfig>,
query: DataQuery,
): CameraManagerEngine | null {
const cameraConfig = cameras.get(query.cameraID);
return cameraConfig ? this.getEngineForCamera(cameraConfig) : null;
return cameraManagerEngine;
}
public getEngineForCamera(cameraConfig?: CameraConfig): CameraManagerEngine | null {
@@ -37,10 +34,34 @@ export class CameraManagerEngineFactory {
return null;
}
let engineKey: string | null = null;
let engine: Engine | null = null;
if (cameraConfig.frigate.camera_name) {
engineKey = 'frigate';
engine = Engine.Frigate;
}
return engineKey ? this._getOrCreateEngine(engineKey) : null;
return engine ? this.getEngine(engine) : null;
}
public getEnginesForCameraIDs(
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
): CameraManagerEngineCameraIDMap | null {
const output: CameraManagerEngineCameraIDMap = new Map();
for (const cameraID of cameraIDs) {
const cameraConfig = cameras.get(cameraID);
if (!cameraConfig) {
continue;
}
const engine = this.getEngineForCamera(cameraConfig);
if (!engine) {
continue;
}
if (!output.has(engine)) {
output.set(engine, new Set());
}
output.get(engine)?.add(cameraID);
}
return output;
}
}
+21 -20
View File
@@ -1,65 +1,67 @@
import { HomeAssistant } from 'custom-card-helpers';
import { CameraConfig } from '../types';
import { MediaQueriesResults } from "../view/media-queries-results";
import { ViewMedia } from '../view/media';
import {
DataQuery,
EventQuery,
EventQueryResultsMap,
PartialEventQuery,
PartialRecordingQuery,
PartialRecordingSegmentsQuery,
QueryReturnType,
RecordingQuery,
RecordingQueryResultsMap,
RecordingSegmentsQuery,
RecordingSegmentsQueryResultsMap,
} from './types';
import { MediaQueries } from '../view/media-queries';
export const CAMERA_MANAGER_ENGINE_EVENT_LIMIT_DEFAULT = 10000;
export interface CameraManagerEngine {
generateDefaultEventQuery(
cameraID: string,
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialEventQuery,
): EventQuery | null;
): EventQuery[] | null;
generateDefaultRecordingQuery(
cameraID: string,
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialRecordingQuery,
): RecordingQuery | null;
): RecordingQuery[] | null;
generateDefaultRecordingSegmentsQuery(
cameraID: string,
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialRecordingSegmentsQuery,
): RecordingSegmentsQuery | null;
): RecordingSegmentsQuery[] | null;
getEvents(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: EventQuery,
): Promise<QueryReturnType<EventQuery> | null>;
): Promise<EventQueryResultsMap | null>;
getRecordings(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: RecordingQuery,
): Promise<QueryReturnType<RecordingQuery> | null>;
): Promise<RecordingQueryResultsMap | null>;
getRecordingSegments(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: RecordingSegmentsQuery,
): Promise<QueryReturnType<RecordingSegmentsQuery> | null>;
): Promise<RecordingSegmentsQueryResultsMap | null>;
generateMediaFromEvents(
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
query: EventQuery,
results: QueryReturnType<EventQuery>,
): ViewMedia[] | null;
generateMediaFromRecordings(
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
query: RecordingQuery,
results: QueryReturnType<RecordingQuery>,
): ViewMedia[] | null;
@@ -73,10 +75,9 @@ export interface CameraManagerEngine {
favorite: boolean,
): Promise<void>;
areMediaQueriesResultsFresh(
queries: MediaQueries,
results: MediaQueriesResults,
): boolean;
getQueryResultMaxAge(
query: DataQuery
): number | null;
getMediaSeekTime(
hass: HomeAssistant,
+364 -182
View File
@@ -4,18 +4,19 @@ import endOfHour from 'date-fns/endOfHour';
import startOfHour from 'date-fns/startOfHour';
import { CAMERA_BIRDSEYE } from '../../const';
import { CameraConfig, RecordingSegment } from '../../types';
import { MediaQueriesResults } from '../../view/media-queries-results';
import { MediaQueriesClassifier } from '../../view/media-queries-classifier';
import { ViewMedia } from '../../view/media';
import { RecordingSegmentsCache } from '../cache';
import { RequestCache, RecordingSegmentsCache } from '../cache';
import {
CameraManagerEngine,
CAMERA_MANAGER_ENGINE_EVENT_LIMIT_DEFAULT,
} from '../engine';
import { DateRange } from '../range';
import {
DataQuery,
Engine,
EventQuery,
EventQueryResults,
EventQueryResultsMap,
FrigateEventQueryResults,
FrigateRecordingQueryResults,
FrigateRecordingSegmentsQueryResults,
@@ -27,7 +28,10 @@ import {
QueryReturnType,
QueryType,
RecordingQuery,
RecordingQueryResults,
RecordingQueryResultsMap,
RecordingSegmentsQuery,
RecordingSegmentsQueryResultsMap,
} from '../types';
import { FrigateRecording } from './types';
import {
@@ -38,7 +42,6 @@ import {
NativeFrigateRecordingSegmentsQuery,
retainEvent,
} from './requests';
import { MediaQueries } from '../../view/media-queries';
import orderBy from 'lodash-es/orderBy';
import throttle from 'lodash-es/throttle';
import { runWhenIdleIfSupported } from '../../utils/basic';
@@ -78,6 +81,7 @@ class FrigateQueryResultsClassifier {
export class FrigateCameraManagerEngine implements CameraManagerEngine {
protected _recordingSegmentsCache: RecordingSegmentsCache;
protected _requestCache: RequestCache;
// Garbage collect segments at most once an hour.
protected _throttledSegmentGarbageCollector = throttle(
@@ -86,8 +90,12 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
{ leading: false, trailing: true },
);
constructor(recordingSegmentsCache: RecordingSegmentsCache) {
constructor(
recordingSegmentsCache: RecordingSegmentsCache,
requestCache: RequestCache,
) {
this._recordingSegmentsCache = recordingSegmentsCache;
this._requestCache = requestCache;
}
public getMediaDownloadPath(
@@ -113,46 +121,81 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
}
public generateDefaultEventQuery(
cameraID: string,
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialEventQuery,
): EventQuery | null {
return {
type: QueryType.Event,
cameraID: cameraID,
...(cameraConfig.frigate.label && { what: cameraConfig.frigate.label }),
...(cameraConfig.frigate.zone && { where: cameraConfig.frigate.zone }),
...query,
};
): EventQuery[] | null {
const relevantCameraConfigs = Array.from(cameraIDs).map((cameraID) =>
cameras.get(cameraID),
);
// If there isn't a label or zone specified, we can come up with a single
// batch query for Frigate that will match across all cameras.
const canDoBatchQuery = relevantCameraConfigs.every(
(cameraConfig) => !cameraConfig?.frigate.label && !cameraConfig?.frigate.zone,
);
if (canDoBatchQuery) {
return [
{
type: QueryType.Event,
cameraIDs: cameraIDs,
...query,
},
];
}
const output: EventQuery[] = [];
for (const cameraID of cameraIDs) {
const cameraConfig = cameras.get(cameraID);
if (cameraConfig) {
output.push({
type: QueryType.Event,
cameraIDs: new Set([cameraID]),
...(cameraConfig.frigate.label && {
what: new Set([cameraConfig.frigate.label]),
}),
...(cameraConfig.frigate.zone && {
where: new Set([cameraConfig.frigate.zone]),
}),
...query,
});
}
}
return output.length ? output : null;
}
public generateDefaultRecordingQuery(
cameraID: string,
_cameraConfig: CameraConfig,
_cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialRecordingQuery,
): RecordingQuery | null {
return {
type: QueryType.Recording,
cameraID: cameraID,
...query,
};
): RecordingQuery[] | null {
return [
{
type: QueryType.Recording,
cameraIDs: cameraIDs,
...query,
},
];
}
public generateDefaultRecordingSegmentsQuery(
cameraID: string,
_cameraConfig: CameraConfig,
_cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
query: PartialRecordingSegmentsQuery,
): RecordingSegmentsQuery | null {
): RecordingSegmentsQuery[] | null {
if (!query.start || !query.end) {
return null;
}
return {
type: QueryType.RecordingSegments,
cameraID: cameraID,
start: query.start,
end: query.end,
...query,
};
return [
{
type: QueryType.RecordingSegments,
cameraIDs: cameraIDs,
start: query.start,
end: query.end,
...query,
},
];
}
public async favoriteMedia(
@@ -169,150 +212,275 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
media.setFavorite(favorite);
}
protected _buildInstanceToCameraIDMapFromQuery(
cameras: Map<string, CameraConfig>,
query: DataQuery,
): Map<string, Set<string>> {
const output: Map<string, Set<string>> = new Map();
for (const cameraID of query.cameraIDs) {
const cameraConfig = this._getQueryableCameraConfig(cameras, cameraID);
const clientID = cameraConfig?.frigate.client_id;
if (clientID) {
if (!output.has(clientID)) {
output.set(clientID, new Set());
}
output.get(clientID)?.add(cameraID);
}
}
return output;
}
protected _getFrigateCameraNamesForCameraIDs(
cameras: Map<string, CameraConfig>,
cameraIDs: Set<string>,
): Set<string> {
const output = new Set<string>();
for (const cameraID of cameraIDs) {
const cameraConfig = this._getQueryableCameraConfig(cameras, cameraID);
if (cameraConfig?.frigate.camera_name) {
output.add(cameraConfig.frigate.camera_name);
}
}
return output;
}
public async getEvents(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: EventQuery,
): Promise<QueryReturnType<EventQuery> | null> {
const cameraConfig = this._getQueryableCameraConfig(cameras, query.cameraID);
if (!cameraConfig) {
return null;
}
): Promise<EventQueryResultsMap | null> {
const output: EventQueryResultsMap = new Map();
const nativeQuery: NativeFrigateEventQuery = {
instance_id: cameraConfig.frigate.client_id,
camera: cameraConfig.frigate.camera_name,
...(query.what && { label: query.what }),
...(query.where && { zone: query.where }),
...(query?.end && { before: Math.floor(query.end.getTime() / 1000) }),
...(query?.start && { after: Math.floor(query.start.getTime() / 1000) }),
...(query?.limit && { limit: query.limit }),
...(query?.hasClip && { has_clip: query.hasClip }),
...(query?.hasSnapshot && { has_snapshot: query.hasSnapshot }),
limit: query?.limit ?? CAMERA_MANAGER_ENGINE_EVENT_LIMIT_DEFAULT,
const processInstanceQuery = async (
instanceID: string,
cameraIDs?: Set<string>,
): Promise<void> => {
if (!cameraIDs || !cameraIDs.size) {
return;
}
const instanceQuery = { ...query, cameraIDs: cameraIDs };
const cachedResult = this._requestCache.get(instanceQuery);
if (cachedResult) {
output.set(query, cachedResult as EventQueryResults);
return;
}
const nativeQuery: NativeFrigateEventQuery = {
instance_id: instanceID,
cameras: Array.from(this._getFrigateCameraNamesForCameraIDs(cameras, cameraIDs)),
...(query.what && { label: Array.from(query.what) }),
...(query.where && { zone: Array.from(query.where) }),
...(query.end && { before: Math.floor(query.end.getTime() / 1000) }),
...(query.start && { after: Math.floor(query.start.getTime() / 1000) }),
...(query.limit && { limit: query.limit }),
...(query.hasClip && { has_clip: query.hasClip }),
...(query.hasSnapshot && { has_snapshot: query.hasSnapshot }),
limit: query?.limit ?? CAMERA_MANAGER_ENGINE_EVENT_LIMIT_DEFAULT,
};
const result: FrigateEventQueryResults = {
type: QueryResultsType.Event,
engine: Engine.Frigate,
instanceID: instanceID,
events: await getEvents(hass, nativeQuery),
expiry: add(new Date(), { seconds: EVENT_REQUEST_CACHE_MAX_AGE_SECONDS }),
cached: false,
};
this._requestCache.set(query, { ...result, cached: true }, result.expiry);
output.set(instanceQuery, result);
};
return <FrigateEventQueryResults>{
type: QueryResultsType.Event,
engine: Engine.Frigate,
events: await getEvents(hass, nativeQuery),
expiry: add(new Date(), { seconds: EVENT_REQUEST_CACHE_MAX_AGE_SECONDS }),
cached: false,
};
// Frigate allows multiple cameras to be searched for events in a single
// query. Break them down into groups of cameras per Frigate instance, then
// query once per instance for all cameras in that instance.
const instances = this._buildInstanceToCameraIDMapFromQuery(cameras, query);
await Promise.all(
Array.from(instances.keys()).map((instanceID) =>
processInstanceQuery(instanceID, instances.get(instanceID)),
),
);
return output.size ? output : null;
}
public async getRecordings(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: RecordingQuery,
): Promise<QueryReturnType<RecordingQuery> | null> {
const cameraConfig = this._getQueryableCameraConfig(cameras, query.cameraID);
if (!cameraConfig) {
return null;
}
if (!cameraConfig || !cameraConfig.frigate.camera_name) {
return null;
}
): Promise<RecordingQueryResultsMap | null> {
const output: RecordingQueryResultsMap = new Map();
const recordingSummary = await getRecordingsSummary(
hass,
cameraConfig.frigate.client_id,
cameraConfig.frigate.camera_name,
);
const processQuery = async (query: RecordingQuery): Promise<void> => {
const cachedResult = this._requestCache.get(query);
if (cachedResult) {
output.set(query, cachedResult as RecordingQueryResults);
return;
}
let recordings: FrigateRecording[] = [];
// There will only ever be a single cameraID specified for queries in this
// inner function.
const cameraID = [...query.cameraIDs][0];
const cameraConfig = this._getQueryableCameraConfig(cameras, cameraID);
if (!cameraConfig || !cameraConfig.frigate.camera_name) {
return;
}
for (const dayData of recordingSummary ?? []) {
for (const hourData of dayData.hours) {
const hour = add(dayData.day, { hours: hourData.hour });
const startHour = startOfHour(hour);
const endHour = endOfHour(hour);
if (
(!query.start || startHour >= query.start) &&
(!query.end || endHour <= query.end)
) {
recordings.push({
cameraID: query.cameraID,
startTime: startHour,
endTime: endHour,
events: hourData.events,
});
const recordingSummary = await getRecordingsSummary(
hass,
cameraConfig.frigate.client_id,
cameraConfig.frigate.camera_name,
);
let recordings: FrigateRecording[] = [];
for (const dayData of recordingSummary ?? []) {
for (const hourData of dayData.hours) {
const hour = add(dayData.day, { hours: hourData.hour });
const startHour = startOfHour(hour);
const endHour = endOfHour(hour);
if (
(!query.start || startHour >= query.start) &&
(!query.end || endHour <= query.end)
) {
recordings.push({
cameraID: cameraID,
startTime: startHour,
endTime: endHour,
events: hourData.events,
});
}
}
}
}
if (query.limit !== undefined) {
// Frigate does not natively support a way to limit recording searches so
// this simulates it.
recordings = orderBy(
recordings,
(recording: FrigateRecording) => recording.startTime,
'desc',
).slice(0, query.limit);
}
if (query.limit !== undefined) {
// Frigate does not natively support a way to limit recording searches so
// this simulates it.
recordings = orderBy(
recordings,
(recording: FrigateRecording) => recording.startTime,
'desc',
).slice(0, query.limit);
}
return <FrigateRecordingQueryResults>{
type: QueryResultsType.Recording,
engine: Engine.Frigate,
recordings: recordings,
expiry: add(new Date(), {
seconds: RECORDING_SUMMARY_REQUEST_CACHE_MAX_AGE_SECONDS,
}),
cached: false,
const result: FrigateRecordingQueryResults = {
type: QueryResultsType.Recording,
engine: Engine.Frigate,
instanceID: cameraConfig.frigate.client_id,
recordings: recordings,
expiry: add(new Date(), {
seconds: RECORDING_SUMMARY_REQUEST_CACHE_MAX_AGE_SECONDS,
}),
cached: false,
};
this._requestCache.set(query, { ...result, cached: true }, result.expiry);
output.set(query, result);
};
// Frigate recordings can only be queried for a single camera, so fan out
// the inbound query into multiple outbound queries.
await Promise.all(
Array.from(query.cameraIDs).map((cameraID) =>
processQuery({ ...query, cameraIDs: new Set([cameraID]) }),
),
);
return output.size ? output : null;
}
public async getRecordingSegments(
hass: HomeAssistant,
cameras: Map<string, CameraConfig>,
query: RecordingSegmentsQuery,
): Promise<QueryReturnType<RecordingSegmentsQuery> | null> {
const cameraConfig = this._getQueryableCameraConfig(cameras, query.cameraID);
if (!cameraConfig || !cameraConfig.frigate.camera_name) {
return null;
}
): Promise<RecordingSegmentsQueryResultsMap | null> {
const output: RecordingSegmentsQueryResultsMap = new Map();
const range: DateRange = { start: query.start, end: query.end };
const processQuery = async (query: RecordingSegmentsQuery): Promise<void> => {
// There will only ever be a single cameraID specified for queries in this
// inner function.
const cameraID = [...query.cameraIDs][0];
const cameraConfig = this._getQueryableCameraConfig(cameras, cameraID);
if (!cameraConfig || !cameraConfig.frigate.camera_name) {
return;
}
// A note on Frigate Recording Segments:
// - Unlike other query types, there is an internal cache at the engine
// level for segments to allow caching "within an existing query" (e.g. if
// we already cached hour 1-8, we will avoid a fetch if we request hours
// 2-3 even though the query is different -- the segments won't be). This
// is since the volume of data in segment transfers can be high, and the
// segments can be used in high frequency situations (e.g. video seeking).
const cachedSegments = this._recordingSegmentsCache.get(query.cameraID, range);
if (cachedSegments) {
return {
const range: DateRange = { start: query.start, end: query.end };
// A note on Frigate Recording Segments:
// - There is an internal cache at the engine level for segments to allow
// caching "within an existing query" (e.g. if we already cached hour
// 1-8, we will avoid a fetch if we request hours 2-3 even though the
// query is different -- the segments won't be). This is since the
// volume of data in segment transfers can be high, and the segments can
// be used in high frequency situations (e.g. video seeking).
const cachedSegments = this._recordingSegmentsCache.get(cameraID, range);
if (cachedSegments) {
output.set(query, <FrigateRecordingSegmentsQueryResults>{
type: QueryResultsType.RecordingSegments,
engine: Engine.Frigate,
instanceID: cameraConfig.frigate.client_id,
segments: cachedSegments,
cached: true,
});
return;
}
const request: NativeFrigateRecordingSegmentsQuery = {
instance_id: cameraConfig.frigate.client_id,
camera: cameraConfig.frigate.camera_name,
after: Math.floor(query.start.getTime() / 1000),
before: Math.floor(query.end.getTime() / 1000),
};
const segments = await getRecordingSegments(hass, request);
this._recordingSegmentsCache.add(cameraID, range, segments);
output.set(query, <FrigateRecordingSegmentsQueryResults>{
type: QueryResultsType.RecordingSegments,
engine: Engine.Frigate,
segments: cachedSegments,
cached: true,
};
}
const request: NativeFrigateRecordingSegmentsQuery = {
instance_id: cameraConfig.frigate.client_id,
camera: cameraConfig.frigate.camera_name,
after: Math.floor(query.start.getTime() / 1000),
before: Math.floor(query.end.getTime() / 1000),
instanceID: cameraConfig.frigate.client_id,
segments: segments,
cached: false,
});
};
const segments = await getRecordingSegments(hass, request);
this._recordingSegmentsCache.add(query.cameraID, range, segments);
// Frigate recording segments can only be queried for a single camera, so
// fan out the inbound query into multiple outbound queries.
await Promise.all(
Array.from(query.cameraIDs).map((cameraID) =>
processQuery({ ...query, cameraIDs: new Set([cameraID]) }),
),
);
runWhenIdleIfSupported(() => this._throttledSegmentGarbageCollector(hass, cameras));
return output.size ? output : null;
}
return {
type: QueryResultsType.RecordingSegments,
engine: Engine.Frigate,
segments: segments,
cached: false,
};
protected _getCameraIDMatch(
cameras: Map<string, CameraConfig>,
query: DataQuery,
instanceID: string,
cameraName: string,
): string | null {
// If the query is only for a single cameraID, all results are assumed to
// belong to it for performance reasons. Otherwise, we need to map the
// instanceID and camera name for the known cameras, and get the precise
// cameraID that matches the expected instance ID / camera name.
if (query.cameraIDs.size === 1) {
return [...query.cameraIDs][0];
}
for (const [cameraID, cameraConfig] of cameras.entries()) {
if (
cameraConfig.frigate.client_id === instanceID &&
cameraConfig.frigate.camera_name === cameraName
) {
return cameraID;
}
}
return null;
}
public generateMediaFromEvents(
cameraConfig: CameraConfig,
cameras: Map<string, CameraConfig>,
query: EventQuery,
results: QueryReturnType<EventQuery>,
): ViewMedia[] | null {
@@ -322,6 +490,19 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
const output: ViewMedia[] = [];
for (const event of results.events) {
const cameraID = this._getCameraIDMatch(
cameras,
query,
results.instanceID,
event.camera,
);
if (!cameraID) {
continue;
}
const cameraConfig = this._getQueryableCameraConfig(cameras, cameraID);
if (!cameraConfig) {
continue;
}
let mediaType: 'clip' | 'snapshot' | null = null;
if (
!query.hasClip &&
@@ -339,7 +520,7 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
}
const media = FrigateViewMediaFactory.createEventViewMedia(
mediaType,
query.cameraID,
cameraID,
cameraConfig,
event,
);
@@ -351,8 +532,8 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
}
public generateMediaFromRecordings(
cameraConfig: CameraConfig,
query: RecordingQuery,
cameras: Map<string, CameraConfig>,
_query: RecordingQuery,
results: QueryReturnType<RecordingQuery>,
): ViewMedia[] | null {
if (!FrigateQueryResultsClassifier.isFrigateRecordingQueryResults(results)) {
@@ -361,8 +542,12 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
const output: ViewMedia[] = [];
for (const recording of results.recordings) {
const cameraConfig = this._getQueryableCameraConfig(cameras, recording.cameraID);
if (!cameraConfig) {
continue;
}
const media = FrigateViewMediaFactory.createRecordingViewMedia(
query.cameraID,
recording.cameraID,
recording,
cameraConfig,
);
@@ -373,23 +558,13 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
return output;
}
public areMediaQueriesResultsFresh(
queries: MediaQueries,
results: MediaQueriesResults,
): boolean {
let freshThreshold: number | null = null;
if (MediaQueriesClassifier.areEventQueries(queries)) {
freshThreshold = EVENT_REQUEST_CACHE_MAX_AGE_SECONDS;
} else if (MediaQueriesClassifier.areRecordingQueries(queries)) {
freshThreshold = RECORDING_SUMMARY_REQUEST_CACHE_MAX_AGE_SECONDS;
public getQueryResultMaxAge(query: DataQuery): number | null {
if (query.type === QueryType.Event) {
return EVENT_REQUEST_CACHE_MAX_AGE_SECONDS;
} else if (query.type === QueryType.Recording) {
return RECORDING_SUMMARY_REQUEST_CACHE_MAX_AGE_SECONDS;
}
const now = new Date();
const resultsTimestamp = results.getResultsTimestamp();
return (
!freshThreshold ||
!resultsTimestamp ||
add(resultsTimestamp, { seconds: freshThreshold }) >= now
);
return null;
}
public async getMediaSeekTime(
@@ -404,18 +579,26 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
return null;
}
const cameraID = media.getCameraID();
const query: RecordingSegmentsQuery = {
cameraID: media.getCameraID(),
cameraIDs: new Set([cameraID]),
start: start,
end: end,
type: QueryType.RecordingSegments,
};
const segments = await this.getRecordingSegments(hass, cameras, query);
const out = segments
? this._getSeekTimeInSegments(start, target, segments.segments)
: null;
return out;
const results = await this.getRecordingSegments(hass, cameras, query);
if (results) {
return this._getSeekTimeInSegments(
start,
target,
// There will only be a single result since Frigate recording segments
// searches are per camera which is specified singularly above.
Array.from(results.values())[0].segments,
);
}
return null;
}
protected _getQueryableCameraConfig(
@@ -438,10 +621,10 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
cameras: Map<string, CameraConfig>,
): Promise<void> {
const cameraIDs = this._recordingSegmentsCache.getCameraIDs();
const recordingQueries: RecordingQuery[] = cameraIDs.map((cameraID) => ({
cameraID: cameraID,
const recordingQuery: RecordingQuery = {
cameraIDs: new Set(cameraIDs),
type: QueryType.Recording,
}));
};
const countSegments = () =>
sum(
@@ -451,19 +634,6 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
);
const segmentsStart = countSegments();
const results: Map<RecordingQuery, FrigateRecordingQueryResults> = new Map();
await Promise.all(
recordingQueries.map((query) =>
(async () => {
const recordings = await this.getRecordings(hass, cameras, query);
if (recordings && recordings.engine === Engine.Frigate) {
results.set(query, recordings as FrigateRecordingQueryResults);
}
})(),
),
);
// Performance: _recordingSegments is potentially very large (e.g. 10K - 1M
// items) and each item must be examined, so care required here to stick to
// nothing worse than O(n) performance.
@@ -471,16 +641,28 @@ export class FrigateCameraManagerEngine implements CameraManagerEngine {
return `${cameraID}/${startTime.getDate()}/${startTime.getHours()}`;
};
const results = await this.getRecordings(hass, cameras, recordingQuery);
if (!results) {
return;
}
for (const [query, result] of results) {
if (!FrigateQueryResultsClassifier.isFrigateRecordingQueryResults(result)) {
continue;
}
const goodHours: Set<string> = new Set();
for (const recording of result.recordings) {
goodHours.add(getHourID(recording.cameraID, recording.startTime));
}
// Frigate recordings are always executed individually, so there'll only
// be a single results.
const cameraID = Array.from(query.cameraIDs)[0];
this._recordingSegmentsCache.expireMatches(
query.cameraID,
cameraID,
(segment: RecordingSegment) => {
const hourID = getHourID(query.cameraID, fromUnixTime(segment.start_time));
const hourID = getHourID(cameraID, fromUnixTime(segment.start_time));
// ~O(1) lookup time for a JS set.
return goodHours.has(hourID);
},
+3 -3
View File
@@ -95,9 +95,9 @@ export async function retainEvent(
export interface NativeFrigateEventQuery {
instance_id?: string;
camera?: string;
label?: string;
zone?: string;
cameras?: string[];
labels?: string[];
zones?: string[];
after?: number;
before?: number;
limit?: number;
+119 -89
View File
@@ -5,6 +5,7 @@ import {
DataQuery,
EventQuery,
EventQueryResults,
EventQueryResultsMap,
MediaQuery,
PartialDataQuery,
PartialEventQuery,
@@ -17,16 +18,21 @@ import {
QueryType,
RecordingQuery,
RecordingQueryResults,
RecordingQueryResultsMap,
RecordingSegmentsQuery,
RecordingSegmentsQueryResults,
RecordingSegmentsQueryResultsMap,
ResultsMap,
} from './types.js';
import orderBy from 'lodash-es/orderBy';
import { CameraManagerEngineFactory } from './engine-factory.js';
import { ViewMedia } from '../view/media.js';
import { MediaQueriesResults } from '../view/media-queries-results';
import { MemoryRequestCache } from './cache.js';
import { MediaQueries } from '../view/media-queries.js';
import uniqBy from 'lodash-es/uniqBy';
import { CameraManagerEngine } from './engine.js';
import sum from 'lodash-es/sum';
import add from 'date-fns/add';
export class QueryClassifier {
public static isEventQuery(query: DataQuery | PartialDataQuery): query is EventQuery {
@@ -67,27 +73,22 @@ export interface ExtendedMediaQueryResult<T extends MediaQuery> {
results: ViewMedia[];
}
export type RequestCache = MemoryRequestCache<DataQuery, QueryResults>;
export class CameraManager {
protected _engineFactory: CameraManagerEngineFactory;
protected _cameras: Map<string, CameraConfig>;
protected _requestCache: RequestCache;
constructor(
engineFactory: CameraManagerEngineFactory,
cameras: Map<string, CameraConfig>,
requestCache: RequestCache,
) {
this._engineFactory = engineFactory;
this._cameras = cameras;
this._requestCache = requestCache;
}
public generateDefaultEventQueries(
cameraIDs: string | Set<string>,
partialQuery: PartialEventQuery,
): EventQuery[] {
): EventQuery[] | null {
return this._generateDefaultQueries(cameraIDs, {
...partialQuery,
type: QueryType.Event,
@@ -97,7 +98,7 @@ export class CameraManager {
public generateDefaultRecordingQueries(
cameraIDs: string | Set<string>,
partialQuery: PartialRecordingQuery,
): RecordingQuery[] {
): RecordingQuery[] | null {
return this._generateDefaultQueries(cameraIDs, {
...partialQuery,
type: QueryType.Recording,
@@ -107,73 +108,76 @@ export class CameraManager {
public generateDefaultRecordingSegmentsQueries(
cameraIDs: string | Set<string>,
partialQuery: PartialRecordingSegmentsQuery,
): RecordingSegmentsQuery[] {
): RecordingSegmentsQuery[] | null {
return this._generateDefaultQueries(cameraIDs, {
...partialQuery,
type: QueryType.RecordingSegments,
});
}
protected _generateDefaultQueries<PQT extends Partial<DataQuery>>(
protected _generateDefaultQueries<PQT extends PartialDataQuery>(
cameraIDs: string | Set<string>,
partialQuery: PQT,
): PartialQueryConcreteType<PQT>[] {
): PartialQueryConcreteType<PQT>[] | null {
const concreteQueries: PartialQueryConcreteType<PQT>[] = [];
const _cameraIDs = setify(cameraIDs);
_cameraIDs.forEach((cameraID) => {
const cameraConfig = this._cameras.get(cameraID);
if (!cameraConfig) {
return;
}
const engines = this._engineFactory.getEnginesForCameraIDs(
this._cameras,
_cameraIDs,
);
const engine = this._engineFactory.getEngineForCamera(cameraConfig);
if (!engine) {
return;
}
if (!engines) {
return null;
}
let query: DataQuery | null = null;
for (const [engine, cameraIDs] of engines) {
let queries: DataQuery[] | null = null;
if (QueryClassifier.isEventQuery(partialQuery)) {
query = engine.generateDefaultEventQuery(cameraID, cameraConfig, partialQuery);
queries = engine.generateDefaultEventQuery(
this._cameras,
cameraIDs,
partialQuery,
);
} else if (QueryClassifier.isRecordingQuery(partialQuery)) {
query = engine.generateDefaultRecordingQuery(
cameraID,
cameraConfig,
queries = engine.generateDefaultRecordingQuery(
this._cameras,
cameraIDs,
partialQuery,
);
} else if (QueryClassifier.isRecordingSegmentsQuery(partialQuery)) {
query = engine.generateDefaultRecordingSegmentsQuery(
cameraID,
cameraConfig,
queries = engine.generateDefaultRecordingSegmentsQuery(
this._cameras,
cameraIDs,
partialQuery,
);
}
if (query) {
for (const query of queries ?? []) {
concreteQueries.push(query as PartialQueryConcreteType<PQT>);
}
});
}
return concreteQueries;
}
public async getEvents(
hass: HomeAssistant,
query: EventQuery | EventQuery[],
): Promise<Map<EventQuery, EventQueryResults>> {
): Promise<EventQueryResultsMap> {
return await this._handleQuery(hass, query);
}
public async getRecordings(
hass: HomeAssistant,
query: RecordingQuery | RecordingQuery[],
): Promise<Map<RecordingQuery, RecordingQueryResults>> {
): Promise<RecordingQueryResultsMap> {
return await this._handleQuery(hass, query);
}
public async getRecordingSegments(
hass: HomeAssistant,
query: RecordingSegmentsQuery | RecordingSegmentsQuery[],
): Promise<Map<RecordingSegmentsQuery, RecordingSegmentsQueryResults>> {
): Promise<RecordingSegmentsQueryResultsMap> {
return await this._handleQuery(hass, query);
}
@@ -301,16 +305,29 @@ export class CameraManager {
queries: MediaQueries,
results: MediaQueriesResults,
): boolean {
const cameraIDs: Set<string> = new Set();
(queries.getQueries() ?? []).forEach((query) => cameraIDs.add(query.cameraID));
for (const cameraID of cameraIDs) {
const cameraConfig = this._cameras.get(cameraID);
if (!cameraConfig) {
return false;
}
const engine = this._engineFactory.getEngineForCamera(cameraConfig);
if (!engine || !engine.areMediaQueriesResultsFresh(queries, results)) {
return false;
const now = new Date();
const resultsTimestamp = results.getResultsTimestamp();
if (!resultsTimestamp) {
return false;
}
for (const query of queries.getQueries() ?? []) {
const engines = this._engineFactory.getEnginesForCameraIDs(
this._cameras,
query.cameraIDs,
);
for (const [engine, cameraIDs] of engines ?? []) {
const maxAgeSeconds = engine.getQueryResultMaxAge({
...query,
cameraIDs: cameraIDs,
});
if (
maxAgeSeconds !== null &&
add(resultsTimestamp, { seconds: maxAgeSeconds }) < now
) {
return false;
}
}
}
return true;
@@ -344,83 +361,96 @@ export class CameraManager {
): Promise<Map<QT, QueryReturnType<QT>>> {
const _queries = arrayify(query);
const results = new Map<QT, QueryReturnType<QT>>();
const queryStartTime = new Date();
let queryCachedCount = 0;
const processEngineQuery = async (
engine: CameraManagerEngine,
query?: QT,
): Promise<void> => {
if (!query) {
return;
}
let engineResult: Map<QT, QueryReturnType<QT>> | null = null;
if (QueryClassifier.isEventQuery(query)) {
engineResult = (await engine.getEvents(hass, this._cameras, query)) as Map<
QT,
QueryReturnType<QT>
> | null;
} else if (QueryClassifier.isRecordingQuery(query)) {
engineResult = (await engine.getRecordings(hass, this._cameras, query)) as Map<
QT,
QueryReturnType<QT>
> | null;
} else if (QueryClassifier.isRecordingSegmentsQuery(query)) {
engineResult = (await engine.getRecordingSegments(
hass,
this._cameras,
query,
)) as Map<QT, QueryReturnType<QT>> | null;
}
engineResult?.forEach((value, key) => results.set(key, value));
};
const processQuery = async (query: QT): Promise<void> => {
const cachedResult: QueryReturnType<QT> | null = this._requestCache.get(
query,
) as QueryReturnType<QT> | null;
if (cachedResult) {
queryCachedCount++;
results.set(query, cachedResult);
const engines = this._engineFactory.getEnginesForCameraIDs(
this._cameras,
query.cameraIDs,
);
if (!engines) {
return;
}
const engine = this._engineFactory.getEngineForQuery(this._cameras, query);
if (!engine) {
return;
}
let result: QueryResults | null = null;
if (QueryClassifier.isEventQuery(query)) {
result = await engine.getEvents(hass, this._cameras, query);
} else if (QueryClassifier.isRecordingQuery(query)) {
result = await engine.getRecordings(hass, this._cameras, query);
} else if (QueryClassifier.isRecordingSegmentsQuery(query)) {
result = await engine.getRecordingSegments(hass, this._cameras, query);
}
// The engine may independently cached the results. Respect that in our
// debug logging.
if (result?.cached) {
queryCachedCount++;
}
if (result) {
if (result.expiry) {
this._requestCache.set(query, { ...result, cached: true }, result.expiry);
}
results.set(query, result as QueryReturnType<QT>);
}
await Promise.all(
Array.from(engines.keys()).map((engine) =>
processEngineQuery(engine, { ...query, cameraIDs: engines.get(engine) }),
),
);
};
await Promise.all(_queries.map((query) => processQuery(query)));
const cachedOutputQueries = sum(
Array.from(results.values()).map((result) => Number(result.cached)),
);
console.debug(
'Frigate Card CameraManager request (Cached:',
`${queryCachedCount}/${_queries.length},`,
`Duration: ${(new Date().getTime() - queryStartTime.getTime()) / 1000}s,`,
'Queries:',
'Frigate Card CameraManager request [Input queries:',
_queries.length,
', Cached output queries:',
cachedOutputQueries,
', Total output queries:',
results.size,
', Duration:',
`${(new Date().getTime() - queryStartTime.getTime()) / 1000}s,`,
', Queries:',
_queries,
', Results:',
results,
')',
']',
);
return results;
}
protected _convertQueryResultsToMedia<QT extends DataQuery>(
results: Map<QT, QueryReturnType<QT>>,
results: ResultsMap<QT>,
): ViewMedia[] {
const mediaArray: ViewMedia[] = [];
for (const [query, result] of results.entries()) {
const cameraConfig = this._cameras.get(query.cameraID);
const engine = this._engineFactory.getEngineForCamera(cameraConfig);
const engine = this._engineFactory.getEngine(result.engine);
if (engine && cameraConfig) {
if (engine) {
let media: ViewMedia[] | null = null;
if (
QueryClassifier.isEventQuery(query) &&
QueryResultClassifier.isEventQueryResult(result)
) {
media = engine.generateMediaFromEvents(cameraConfig, query, result);
media = engine.generateMediaFromEvents(this._cameras, query, result);
} else if (
QueryClassifier.isRecordingQuery(query) &&
QueryResultClassifier.isRecordingQuery(result)
) {
media = engine.generateMediaFromRecordings(cameraConfig, query, result);
media = engine.generateMediaFromRecordings(this._cameras, query, result);
}
if (media) {
mediaArray.push(...media);
+11 -3
View File
@@ -23,7 +23,7 @@ export enum Engine {
export interface DataQuery {
type: QueryType;
cameraID: string;
cameraIDs: Set<string>;
}
export type PartialDataQuery = Partial<DataQuery>;
@@ -63,6 +63,11 @@ export type PartialQueryConcreteType<PQT> = PQT extends PartialEventQuery
? RecordingSegmentsQuery
: never;
export type ResultsMap<QT> = Map<QT, QueryReturnType<QT>>;
export type EventQueryResultsMap = ResultsMap<EventQuery>;
export type RecordingQueryResultsMap = ResultsMap<RecordingQuery>;
export type RecordingSegmentsQueryResultsMap = ResultsMap<RecordingSegmentsQuery>;
// ===========
// Event Query
// ===========
@@ -77,10 +82,10 @@ export interface EventQuery extends MediaQuery {
hasClip?: boolean;
// Frigate equivalent: label
what?: string;
what?: Set<string>;
// Frigate equivalent: zone
where?: string;
where?: Set<string>;
}
export type PartialEventQuery = Partial<EventQuery>;
@@ -121,15 +126,18 @@ export interface RecordingSegmentsQueryResults extends QueryResults {
export interface FrigateEventQueryResults extends EventQueryResults {
engine: Engine.Frigate;
instanceID: string;
events: FrigateEvent[];
}
export interface FrigateRecordingQueryResults extends RecordingQueryResults {
engine: Engine.Frigate;
instanceID: string;
recordings: FrigateRecording[];
}
export interface FrigateRecordingSegmentsQueryResults
extends RecordingSegmentsQueryResults {
engine: Engine.Frigate;
instanceID: string;
}
-1
View File
@@ -1106,7 +1106,6 @@ export class FrigateCard extends LitElement {
this._cameraManager = new CameraManager(
new CameraManagerEngineFactory(),
this._cameras,
new RequestCache(),
);
}
+7 -4
View File
@@ -429,7 +429,7 @@ export class FrigateCardTimelineCore extends LitElement {
}) // Whether or not to set the timeline window.
.mergeInContext({
...(canSeek && { mediaViewer: { seek: targetTime } }),
...this._setWindowInContext(properties)
...this._setWindowInContext(properties),
})
.dispatchChangeEvent(this);
}
@@ -616,9 +616,12 @@ export class FrigateCardTimelineCore extends LitElement {
options?.window ?? this._timeline.getWindow(),
);
return new EventMediaQueries(
this._timelineSource.getTimelineEventQueries(cacheFriendlyWindow),
);
const eventQueries =
this._timelineSource.getTimelineEventQueries(cacheFriendlyWindow);
if (!eventQueries) {
return null;
}
return new EventMediaQueries(eventQueries);
}
protected async _createViewWithEventMediaQuery(
+11 -4
View File
@@ -53,12 +53,15 @@ export const createViewForEvents = async (
? options.cameraIDs
: new Set(getAllDependentCameras(cameras, view.camera));
const queries = cameraManager.generateDefaultEventQueries(cameraIDs, {
const eventQueries = cameraManager.generateDefaultEventQueries(cameraIDs, {
...(options?.limit && { limit: options.limit }),
...(options?.mediaType === 'clips' && { hasClip: true }),
...(options?.mediaType === 'snapshots' && { hasSnapshot: true }),
});
query = new EventMediaQueries(queries);
if (!eventQueries) {
return null;
}
query = new EventMediaQueries(eventQueries);
}
let queryResults: MediaQueriesResults | null;
@@ -137,12 +140,16 @@ export const createViewForRecordings = async (
? options.cameraIDs
: new Set(getAllDependentCameras(cameras, view.camera));
const queries = cameraManager.generateDefaultRecordingQueries(cameraIDs, {
const recordingQueries = cameraManager.generateDefaultRecordingQueries(cameraIDs, {
...(options?.start && { start: options.start }),
...(options?.end && { end: options.end }),
});
const query = new RecordingMediaQueries(queries);
if (!recordingQueries) {
return null;
}
const query = new RecordingMediaQueries(recordingQueries);
let queryResults: MediaQueriesResults | null;
try {
+20 -16
View File
@@ -7,7 +7,7 @@ import { ClipsOrSnapshotsOrAll, RecordingSegment } from '../types';
import { CameraManager } from '../camera/manager';
import { EventQuery } from '../camera/types';
import { capEndDate, convertRangeToCacheFriendlyTimes } from '../camera/util';
import { EventMediaQueries } from "../view/media-queries";
import { EventMediaQueries } from '../view/media-queries';
import { ViewMedia } from '../view/media';
import { compressRanges, ExpiringMemoryRangeSet, MemoryRangeSet } from '../camera/range';
import { errorToConsole, ModifyInterface } from './basic.js';
@@ -85,10 +85,7 @@ export class TimelineDataSource {
}
}
public async refresh(
hass: HomeAssistant,
window: TimelineWindow,
): Promise<void> {
public async refresh(hass: HomeAssistant, window: TimelineWindow): Promise<void> {
try {
await Promise.all([
this._refreshEvents(hass, window),
@@ -110,7 +107,7 @@ export class TimelineDataSource {
});
}
public getTimelineEventQueries(window: TimelineWindow): EventQuery[] {
public getTimelineEventQueries(window: TimelineWindow): EventQuery[] | null {
return this._cameraManager.generateDefaultEventQueries(this._cameraIDs, {
start: window.start,
end: window.end,
@@ -135,9 +132,11 @@ export class TimelineDataSource {
}
const cacheFriendlyWindow = this.getCacheFriendlyEventWindow(window);
const query = new EventMediaQueries(
this.getTimelineEventQueries(cacheFriendlyWindow),
);
const eventQueries = this.getTimelineEventQueries(cacheFriendlyWindow)
if (!eventQueries) {
return;
}
const query = new EventMediaQueries(eventQueries);
const results = await this._cameraManager.executeMediaQueries(hass, query);
for (const media of results?.getResults() ?? []) {
@@ -224,7 +223,7 @@ export class TimelineDataSource {
endCap: true,
});
const queries = this._cameraManager.generateDefaultRecordingSegmentsQueries(
const recordingQueries = this._cameraManager.generateDefaultRecordingSegmentsQueries(
this._cameraIDs,
{
start: cacheFriendlyWindow.start,
@@ -232,16 +231,21 @@ export class TimelineDataSource {
},
);
const results = await this._cameraManager.getRecordingSegments(hass, queries);
if (!recordingQueries) {
return;
}
const results = await this._cameraManager.getRecordingSegments(hass, recordingQueries);
const newSegments: Map<string, RecordingSegment[]> = new Map();
for (const [query, result] of results) {
let destination: RecordingSegment[] | undefined = newSegments.get(query.cameraID);
if (!destination) {
destination = [];
newSegments.set(query.cameraID, destination);
for (const cameraID of query.cameraIDs) {
let destination: RecordingSegment[] | undefined = newSegments.get(cameraID);
if (!destination) {
destination = [];
newSegments.set(cameraID, destination);
}
result.segments.forEach((segment) => destination?.push(segment));
}
result.segments.forEach((segment) => destination?.push(segment));
}
for (const [cameraID, segments] of newSegments.entries()) {