Skip to content
This repository was archived by the owner on May 11, 2026. It is now read-only.
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
126 changes: 80 additions & 46 deletions moose/app/apis/aircraftSpeedAltitudeByType.ts
Original file line number Diff line number Diff line change
@@ -1,79 +1,113 @@
import { Api, getMooseUtils, WebApp } from "@514labs/moose-lib";
import { Api, getMooseUtils, WebApp, sql } from "@514labs/moose-lib";
import { AircraftTrackingProcessed_Table } from "../index";
import { AircraftStatsByCategory_Table } from "../ingest/aircraft_stats_mv";
import express, { Request } from "express";
import cors from "cors";

/**
* Parameters for the aircraft speed and altitude by type API
*/
interface AircraftSpeedAltitudeParams {
/** Optional filter for specific aircraft category */
category?: string;
/** Optional minimum barometric altitude filter */
minAltitude?: number;
/** Optional maximum barometric altitude filter */
maxAltitude?: number;
/** Optional minimum ground speed filter */
minSpeed?: number;
/** Optional maximum ground speed filter */
maxSpeed?: number;
}

function hasActiveFilters(params: AircraftSpeedAltitudeParams): boolean {
return !!(
params.category ||
params.minAltitude ||
params.maxAltitude ||
params.minSpeed ||
params.maxSpeed
);
}

/**
* Constructs the SQL query for aircraft speed and altitude statistics
* Fast path: reads ~10 rows from the AggregatingMergeTree target table
* populated by the MaterializedView. Used when no range filters are active.
*/
const buildAircraftStatsQuery = (sql: any, aircraft_cols: any, params: AircraftSpeedAltitudeParams) => {
const buildMVQuery = (
mvCols: typeof AircraftStatsByCategory_Table.columns,
params: AircraftSpeedAltitudeParams,
) => {
const categoryFilter =
params.category
? sql`WHERE ${mvCols.category} = ${params.category}`
: sql``;

return sql`
SELECT
${aircraft_cols.category} as aircraft_category,
${mvCols.category} AS aircraft_category,
sum(${mvCols.total_records}) AS total_records,
avgMerge(${mvCols.avg_alt_baro}) AS avg_barometric_altitude,
min(${mvCols.min_alt_baro}) AS min_barometric_altitude,
max(${mvCols.max_alt_baro}) AS max_barometric_altitude,
stddevPopMerge(${mvCols.stddev_alt_baro}) AS altitude_stddev,
avgMerge(${mvCols.avg_gs}) AS avg_ground_speed,
min(${mvCols.min_gs}) AS min_ground_speed,
max(${mvCols.max_gs}) AS max_ground_speed,
stddevPopMerge(${mvCols.stddev_gs}) AS speed_stddev,
uniqMerge(${mvCols.unique_aircraft}) AS unique_aircraft_count
FROM ${AircraftStatsByCategory_Table}
${categoryFilter}
GROUP BY ${mvCols.category}
ORDER BY total_records DESC
`;
};

/**
* Slow path: full table scan on the raw table.
* Only used when altitude/speed range filters are active.
*/
const buildFullScanQuery = (
sqlTag: any,
cols: any,
params: AircraftSpeedAltitudeParams,
) => {
return sqlTag`
SELECT
${cols.category} as aircraft_category,
COUNT(*) as total_records,
AVG(${aircraft_cols.alt_baro}) as avg_barometric_altitude,
MIN(${aircraft_cols.alt_baro}) as min_barometric_altitude,
MAX(${aircraft_cols.alt_baro}) as max_barometric_altitude,
STDDEV_POP(${aircraft_cols.alt_baro}) as altitude_stddev,
AVG(${aircraft_cols.gs}) as avg_ground_speed,
MIN(${aircraft_cols.gs}) as min_ground_speed,
MAX(${aircraft_cols.gs}) as max_ground_speed,
STDDEV_POP(${aircraft_cols.gs}) as speed_stddev,
COUNT(DISTINCT ${aircraft_cols.hex}) as unique_aircraft_count
AVG(${cols.alt_baro}) as avg_barometric_altitude,
MIN(${cols.alt_baro}) as min_barometric_altitude,
MAX(${cols.alt_baro}) as max_barometric_altitude,
STDDEV_POP(${cols.alt_baro}) as altitude_stddev,
AVG(${cols.gs}) as avg_ground_speed,
MIN(${cols.gs}) as min_ground_speed,
MAX(${cols.gs}) as max_ground_speed,
STDDEV_POP(${cols.gs}) as speed_stddev,
COUNT(DISTINCT ${cols.hex}) as unique_aircraft_count
FROM ${AircraftTrackingProcessed_Table}
WHERE ${aircraft_cols.alt_baro} > 0
AND ${aircraft_cols.gs} > 0
AND ${aircraft_cols.category} != ''
-- Optional category filter
AND (${params.category || ''} = '' OR ${aircraft_cols.category} = ${params.category || ''})
-- Optional altitude range filters
AND (${params.minAltitude || -999999} = -999999 OR ${aircraft_cols.alt_baro} >= ${params.minAltitude || -999999})
AND (${params.maxAltitude || 999999} = 999999 OR ${aircraft_cols.alt_baro} <= ${params.maxAltitude || 999999})
-- Optional speed range filters
AND (${params.minSpeed || -999999} = -999999 OR ${aircraft_cols.gs} >= ${params.minSpeed || -999999})
AND (${params.maxSpeed || 999999} = 999999 OR ${aircraft_cols.gs} <= ${params.maxSpeed || 999999})
GROUP BY ${aircraft_cols.category}
WHERE ${cols.alt_baro} > 0
AND ${cols.gs} > 0
AND ${cols.category} != ''
AND (${params.category || ''} = '' OR ${cols.category} = ${params.category || ''})
AND (${params.minAltitude || -999999} = -999999 OR ${cols.alt_baro} >= ${params.minAltitude || -999999})
AND (${params.maxAltitude || 999999} = 999999 OR ${cols.alt_baro} <= ${params.maxAltitude || 999999})
AND (${params.minSpeed || -999999} = -999999 OR ${cols.gs} >= ${params.minSpeed || -999999})
AND (${params.maxSpeed || 999999} = 999999 OR ${cols.gs} <= ${params.maxSpeed || 999999})
GROUP BY ${cols.category}
ORDER BY total_records DESC
`;
};


const app = express();
app.use(cors());
app.use(express.json());

/**
* Express API Handler
* API that provides speed and altitude statistics for different aircraft types/categories
* Uses barometric altitude (NOT geometric altitude) and ground speed
*/
app.get("/aircraftSpeedAltitudeByType", async (req: Request<{}, {}, {}, AircraftSpeedAltitudeParams>, res) => {
const { client, sql } = await getMooseUtils();

// Reference the source table object inside the function
const aircraft_cols = AircraftTrackingProcessed_Table?.columns;
const moose = await getMooseUtils();
const params = req.query;

// Execute the query with robust optional parameter handling
const query = buildAircraftStatsQuery(sql, aircraft_cols, params);
const result = await client.query.execute(query);
const query = hasActiveFilters(params)
? buildFullScanQuery(
moose.sql,
AircraftTrackingProcessed_Table?.columns,
params,
)
: buildMVQuery(AircraftStatsByCategory_Table?.columns, params);

const result = await moose.client.query.execute(query);
const data = await result.json();
res.json(data);
});
Expand Down
16 changes: 12 additions & 4 deletions moose/app/index.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,12 @@
export * from "./connectors/fetch_and_ingest_military_aircraft";
export * from './apis/aircraftSpeedAltitudeByType';
export * from './apis/mcp';
export * from './ingest/ingest';
// export * from "./connectors/fetch_and_ingest_military_aircraft";
export * from "./apis/aircraftSpeedAltitudeByType";
export * from "./apis/mcp";
export * from "./ingest/ingest";
export * from "./ingest/aircraft_stats_mv";

import { OlapTable } from "@514labs/moose-lib";
import { AircraftTrackingData } from "./datamodels/models";

const t1 = new OlapTable<AircraftTrackingData>("t", {
version: "0.0",
});
68 changes: 68 additions & 0 deletions moose/app/ingest/aircraft_stats_mv.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
import {
OlapTable,
MaterializedView,
ClickHouseEngines,
Aggregated,
SimpleAggregated,
sql,
} from "@514labs/moose-lib";
import { AircraftTrackingProcessed_Table } from "./ingest";

/**
* Pre-aggregated aircraft statistics by category.
*
* Populated incrementally via a MaterializedView on insert into
* AircraftTrackingProcessedTable, so the dashboard no longer needs
* to full-scan 63 M+ rows on every page load.
*/
interface AircraftStatsByCategory {
category: string;
total_records: number & SimpleAggregated<"sum", number>;
sum_alt_baro: number & SimpleAggregated<"sum", number>;
min_alt_baro: number & SimpleAggregated<"min", number>;
max_alt_baro: number & SimpleAggregated<"max", number>;
sum_gs: number & SimpleAggregated<"sum", number>;
min_gs: number & SimpleAggregated<"min", number>;
max_gs: number & SimpleAggregated<"max", number>;
avg_alt_baro: number & Aggregated<"avg", [number]>;
stddev_alt_baro: number & Aggregated<"stddevPop", [number]>;
avg_gs: number & Aggregated<"avg", [number]>;
stddev_gs: number & Aggregated<"stddevPop", [number]>;
unique_aircraft: number & Aggregated<"uniq", [string]>;
}

export const AircraftStatsByCategory_Table =
new OlapTable<AircraftStatsByCategory>("AircraftStatsByCategoryTable", {
engine: ClickHouseEngines.AggregatingMergeTree,
orderByFields: ["category"],
});

const src = AircraftTrackingProcessed_Table;

export const AircraftStatsByCategory_MV =
new MaterializedView<AircraftStatsByCategory>({
materializedViewName: "AircraftStatsByCategoryMV",
selectTables: [src],
targetTable: AircraftStatsByCategory_Table,
selectStatement: sql.statement`
SELECT
${src.columns.category} AS category,
count() AS total_records,
sum(${src.columns.alt_baro}) AS sum_alt_baro,
min(${src.columns.alt_baro}) AS min_alt_baro,
max(${src.columns.alt_baro}) AS max_alt_baro,
sum(${src.columns.gs}) AS sum_gs,
min(${src.columns.gs}) AS min_gs,
max(${src.columns.gs}) AS max_gs,
avgState(${src.columns.alt_baro}) AS avg_alt_baro,
stddevPopState(${src.columns.alt_baro}) AS stddev_alt_baro,
avgState(${src.columns.gs}) AS avg_gs,
stddevPopState(${src.columns.gs}) AS stddev_gs,
uniqState(${src.columns.hex}) AS unique_aircraft
FROM ${src}
WHERE ${src.columns.alt_baro} > 0
AND ${src.columns.gs} > 0
AND ${src.columns.category} != ''
GROUP BY ${src.columns.category}
`,
});
46 changes: 32 additions & 14 deletions moose/app/ingest/ingest.ts
Original file line number Diff line number Diff line change
@@ -1,25 +1,43 @@
import { OlapTable, Stream, IngestApi, DeadLetterQueue } from "@514labs/moose-lib";
import { AircraftTrackingData, AircraftTrackingProcessed } from "../datamodels/models";
import {
OlapTable,
Stream,
IngestApi,
DeadLetterQueue,
} from "@514labs/moose-lib";
import {
AircraftTrackingData,
AircraftTrackingProcessed,
} from "../datamodels/models";
import { transformAircraft } from "../functions/process_aircraft";

//Raw data ingest pipeline
export const AircraftTrackingData_Table = new OlapTable<AircraftTrackingData>("AircraftTrackingDataTable");
export const AircraftTrackingData_Table = new OlapTable<AircraftTrackingData>(
"AircraftTrackingDataTable",
);

export const AircraftTrackingData_Stream = new Stream<AircraftTrackingData>("AircraftTrackingDataStream", {
destination: AircraftTrackingData_Table
});
export const AircraftTrackingData_Stream = new Stream<AircraftTrackingData>(
"AircraftTrackingDataStream",
{
destination: AircraftTrackingData_Table,
},
);

export const AircraftTrackingData_IngestAPI = new IngestApi<AircraftTrackingData>("AircraftTrackingDataIngestAPI", {
destination: AircraftTrackingData_Stream,
deadLetterQueue: new DeadLetterQueue<AircraftTrackingData>("AircraftTrackingDataDLQ")
});
export const AircraftTrackingData_IngestAPI =
new IngestApi<AircraftTrackingData>("AircraftTrackingDataIngestAPI", {
destination: AircraftTrackingData_Stream,
deadLetterQueue: new DeadLetterQueue<AircraftTrackingData>(
"AircraftTrackingDataDLQ",
),
});

//Derivative data model pipeline
export const AircraftTrackingProcessed_Table = new OlapTable<AircraftTrackingProcessed>("AircraftTrackingProcessedTable");
export const AircraftTrackingProcessed_Table =
new OlapTable<AircraftTrackingProcessed>("AircraftTrackingProcessedTable");

export const AircraftTrackingProcessed_Stream = new Stream<AircraftTrackingProcessed>("AircraftTrackingProcessedStream", {
destination: AircraftTrackingProcessed_Table
});
export const AircraftTrackingProcessed_Stream =
new Stream<AircraftTrackingProcessed>("AircraftTrackingProcessedStream", {
destination: AircraftTrackingProcessed_Table,
});

AircraftTrackingData_Stream!.addTransform(
AircraftTrackingProcessed_Stream!,
Expand Down
Loading
Loading