2026-04-14-8 auth and commercial
This commit is contained in:
+301
-248
@@ -15,309 +15,360 @@ var TelemetryAnalyzerService_1;
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.TelemetryAnalyzerService = void 0;
|
||||
const common_1 = require("@nestjs/common");
|
||||
const schedule_1 = require("@nestjs/schedule");
|
||||
const typeorm_1 = require("@nestjs/typeorm");
|
||||
const typeorm_2 = require("typeorm");
|
||||
const schedule_1 = require("@nestjs/schedule");
|
||||
const telemetry_entity_1 = require("./telemetry.entity");
|
||||
const road_stat_entity_1 = require("../maps/road-stat.entity");
|
||||
const candidate_road_entity_1 = require("../maps/candidate-road.entity");
|
||||
const redis_service_1 = require("../common/redis.service");
|
||||
const external_telemetry_service_1 = require("./external-telemetry.service");
|
||||
const redis_service_1 = require("../common/redis.service");
|
||||
const telegram_service_1 = require("../common/telegram.service");
|
||||
const traffic_grid_service_1 = require("../maps/traffic-grid.service");
|
||||
let TelemetryAnalyzerService = TelemetryAnalyzerService_1 = class TelemetryAnalyzerService {
|
||||
telemetryRepo;
|
||||
roadStatRepo;
|
||||
candidateRoadRepo;
|
||||
dataSource;
|
||||
redisService;
|
||||
externalTelemetry;
|
||||
redisService;
|
||||
telegramService;
|
||||
trafficGrid;
|
||||
dataSource;
|
||||
logger = new common_1.Logger(TelemetryAnalyzerService_1.name);
|
||||
TRAFFIC_CACHE_KEY = 'traffic_snapshot';
|
||||
constructor(telemetryRepo, roadStatRepo, candidateRoadRepo, dataSource, redisService, externalTelemetry) {
|
||||
TRAFFIC_CACHE_KEY = 'live_traffic_congested';
|
||||
constructor(telemetryRepo, roadStatRepo, candidateRoadRepo, externalTelemetry, redisService, telegramService, trafficGrid, dataSource) {
|
||||
this.telemetryRepo = telemetryRepo;
|
||||
this.roadStatRepo = roadStatRepo;
|
||||
this.candidateRoadRepo = candidateRoadRepo;
|
||||
this.dataSource = dataSource;
|
||||
this.redisService = redisService;
|
||||
this.externalTelemetry = externalTelemetry;
|
||||
this.redisService = redisService;
|
||||
this.telegramService = telegramService;
|
||||
this.trafficGrid = trafficGrid;
|
||||
this.dataSource = dataSource;
|
||||
}
|
||||
async handleNightlyIntelligence() {
|
||||
this.logger.log('⏰ Starting automated 3 AM intelligence process...');
|
||||
async runDeepIntelligence(hours = 48) {
|
||||
this.logger.log(`🚀 Starting Full Intelligence Pipeline (Window: ${hours}h)...`);
|
||||
await this.syncExternalData(Math.ceil(hours / 24));
|
||||
const speedResult = await this.analyzeRoadSpeeds(hours);
|
||||
const temporalResult = await this.analyzeTimeProfiles(hours);
|
||||
const discoveryResult = await this.discoverNewRoads(hours);
|
||||
await this.refreshTrafficCache();
|
||||
await this.trafficGrid.refreshGrid();
|
||||
const summary = await this.getAnalysisSummary();
|
||||
await this.telegramService.sendIntelligenceReport({
|
||||
syncResult: speedResult.totalPointsProcessed || 0,
|
||||
updatedSegments: speedResult.segmentsUpdated || 0,
|
||||
discoveredRoads: discoveryResult.candidatesFound || 0,
|
||||
timeProfiles: temporalResult.bucketsUpdated || 0,
|
||||
totalPoints: summary.telemetry.total,
|
||||
days: Math.ceil(hours / 24),
|
||||
});
|
||||
this.logger.log('🏁 Intelligence Pipeline Finished Successfully.');
|
||||
return {
|
||||
success: true,
|
||||
timestamp: new Date().toISOString(),
|
||||
speedAnalysis: speedResult,
|
||||
temporalAnalysis: temporalResult,
|
||||
roadDiscovery: discoveryResult,
|
||||
};
|
||||
}
|
||||
handleDeepIntelligence() {
|
||||
this.runDeepIntelligence(240);
|
||||
}
|
||||
async syncExternalData(days = 7) {
|
||||
this.logger.log(`📡 Starting batch sync (Window: ${days} days)...`);
|
||||
try {
|
||||
await this.syncExternalData(1);
|
||||
await this.analyzeRoadSpeeds(24);
|
||||
await this.discoverNewRoads(168);
|
||||
await this.refreshTrafficCache();
|
||||
this.logger.log('✅ 3 AM intelligence process complete.');
|
||||
const tracks = await this.externalTelemetry.fetchCarTracks(days);
|
||||
if (tracks.length === 0)
|
||||
return;
|
||||
this.logger.log(`📥 Saving ${tracks.length} points to database...`);
|
||||
const values = tracks
|
||||
.filter(t => t.driver_id && t.longitude && t.latitude)
|
||||
.map(t => `('${t.driver_id}', ST_SetSRID(ST_Point(${t.longitude}, ${t.latitude}), 4326), ${t.latitude}, ${t.longitude}, ${t.speed || 0}, ${t.heading || 0}, '${t.created_at}')`).join(',');
|
||||
await this.dataSource.query(`
|
||||
INSERT INTO telemetry_logs ("driverId", location, latitude, longitude, speed, heading, timestamp)
|
||||
VALUES ${values}
|
||||
ON CONFLICT DO NOTHING
|
||||
`);
|
||||
this.logger.log(`✅ Batch sync complete.`);
|
||||
}
|
||||
catch (error) {
|
||||
this.logger.error('❌ Nightly intelligence failed:', error.stack);
|
||||
this.logger.error(`❌ Sync failed: ${error.message}`);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
async syncExternalData(days = 1) {
|
||||
this.logger.log(`📡 Starting batch sync (Window: ${days} days)...`);
|
||||
const tracks = await this.externalTelemetry.fetchCarTracks(days);
|
||||
if (!tracks || tracks.length === 0) {
|
||||
this.logger.warn(`⚠️ No tracks found on external server for the last ${days} days.`);
|
||||
return { imported: 0 };
|
||||
}
|
||||
this.logger.log(`✅ Received ${tracks.length} tracks. Commencing batch insertion...`);
|
||||
this.logger.log(`📥 Saving ${tracks.length} points to database...`);
|
||||
const entities = tracks.map(t => ({
|
||||
driverId: t.driver_id,
|
||||
latitude: t.latitude,
|
||||
longitude: t.longitude,
|
||||
speed: t.speed,
|
||||
heading: t.heading,
|
||||
timestamp: new Date(t.created_at || t.timestamp),
|
||||
location: {
|
||||
type: 'Point',
|
||||
coordinates: [t.longitude, t.latitude],
|
||||
},
|
||||
}));
|
||||
const CHUNK_SIZE = 1000;
|
||||
for (let i = 0; i < entities.length; i += CHUNK_SIZE) {
|
||||
const chunk = entities.slice(i, i + CHUNK_SIZE);
|
||||
const logs = this.telemetryRepo.create(chunk);
|
||||
await this.telemetryRepo.save(logs);
|
||||
}
|
||||
return { imported: tracks.length };
|
||||
}
|
||||
async refreshTrafficCache() {
|
||||
this.logger.log('🚀 Refreshing Redis traffic snapshot...');
|
||||
this.logger.log('🚀 Refreshing Redis traffic snapshot (v2.5.2 Optimized)...');
|
||||
const congested = await this.roadStatRepo.query(`
|
||||
SELECT
|
||||
"segmentId",
|
||||
"congestionFactor",
|
||||
ST_AsGeoJSON(geometry) as geojson
|
||||
SELECT "segmentId", "congestionFactor", ST_AsGeoJSON(geometry, 5) as geojson
|
||||
FROM road_segment_stats
|
||||
WHERE "congestionFactor" > 1.1
|
||||
ORDER BY "congestionFactor" DESC
|
||||
LIMIT 2000
|
||||
`);
|
||||
if (congested.length === 0) {
|
||||
await this.redisService.del(this.TRAFFIC_CACHE_KEY);
|
||||
return { cachedCount: 0 };
|
||||
}
|
||||
await this.redisService.set(this.TRAFFIC_CACHE_KEY, congested);
|
||||
this.logger.log(`✅ Cached ${congested.length} congested segments in Redis.`);
|
||||
const snapshot = congested.map(row => ({
|
||||
sid: row.segmentId,
|
||||
cf: parseFloat(row.congestionFactor),
|
||||
geo: JSON.parse(row.geojson)
|
||||
}));
|
||||
const sampleSize = Math.min(congested.length, 3);
|
||||
const topSegments = congested.slice(0, sampleSize).map(s => `${s.segmentId} (F: ${s.congestionFactor})`).join(', ');
|
||||
const totalSizeKB = Math.round(JSON.stringify(snapshot).length / 1024);
|
||||
this.logger.log(`📊 Traffic Snapshot Sample (Top ${sampleSize}): ${topSegments}`);
|
||||
this.logger.log(`📦 Redis Payload Size: ~${totalSizeKB} KB`);
|
||||
await this.redisService.set(this.TRAFFIC_CACHE_KEY, snapshot);
|
||||
this.logger.log(`✅ Redis traffic snapshot updated with ${congested.length} segments.`);
|
||||
return { cachedCount: congested.length };
|
||||
}
|
||||
async analyzeRoadSpeeds(sinceHours = 24) {
|
||||
this.logger.log(`🔍 Starting road speed analysis for last ${sinceHours}h...`);
|
||||
const matchedData = await this.dataSource.query(`
|
||||
WITH matched_points AS (
|
||||
SELECT
|
||||
t.id AS telemetry_id,
|
||||
t.speed,
|
||||
t."driverId",
|
||||
l.osm_id,
|
||||
l.name,
|
||||
l.highway,
|
||||
l.way_4326,
|
||||
ST_Distance(t.location::geography, l.way_4326::geography) AS distance_m
|
||||
FROM telemetry_logs t
|
||||
CROSS JOIN LATERAL (
|
||||
SELECT osm_id, name, highway, ST_Transform(way, 4326) AS way_4326
|
||||
FROM planet_osm_line
|
||||
WHERE highway IS NOT NULL
|
||||
ORDER BY way <-> ST_Transform(t.location::geometry, 3857)
|
||||
LIMIT 1
|
||||
) l
|
||||
WHERE t.timestamp >= NOW() - INTERVAL '${sinceHours} hours'
|
||||
AND t.speed > 2 -- Ignore stationary points / تجاهل النقاط الثابتة
|
||||
AND ST_Distance(t.location::geography, l.way_4326::geography) < 15 -- 15m snap threshold
|
||||
this.logger.log(`🔍 Starting road speed analysis v2.5 for last ${sinceHours}h...`);
|
||||
const query = `
|
||||
WITH grid_points AS (
|
||||
-- Group by 5m grid cell first to reduce spatial join volume (v2.5.1 WoW Performance)
|
||||
SELECT ST_SnapToGrid(ST_Transform(location::geometry, 3857), 5) AS loc,
|
||||
AVG(speed) as speed
|
||||
FROM telemetry_logs
|
||||
WHERE timestamp >= NOW() - INTERVAL '${sinceHours} hours' AND speed > 2
|
||||
GROUP BY loc
|
||||
),
|
||||
matches AS (
|
||||
-- Bulk Spatial Join (20m radius) with GIST optimization
|
||||
SELECT l.osm_id::text as id, l.name, l.highway, g.speed, l.way
|
||||
FROM grid_points g
|
||||
INNER JOIN planet_osm_line l ON
|
||||
l.highway IS NOT NULL AND
|
||||
l.way && ST_Expand(g.loc, 20) AND
|
||||
ST_DWithin(l.way, g.loc, 20)
|
||||
),
|
||||
stats AS (
|
||||
-- Calculate median speed and aggregate samples
|
||||
SELECT id, name, highway,
|
||||
PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY speed) as avg_speed,
|
||||
COUNT(*) as samples,
|
||||
ST_Transform(MIN(way), 4326) as way4326
|
||||
FROM matches
|
||||
GROUP BY id, name, highway
|
||||
HAVING COUNT(*) >= 3
|
||||
)
|
||||
SELECT
|
||||
osm_id::text AS segment_id,
|
||||
name,
|
||||
highway,
|
||||
AVG(speed) AS avg_speed,
|
||||
PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY speed) AS median_speed,
|
||||
COUNT(*) AS sample_count,
|
||||
COUNT(DISTINCT "driverId") AS unique_drivers,
|
||||
ST_AsGeoJSON(MIN(way_4326)) AS geojson
|
||||
FROM matched_points
|
||||
GROUP BY osm_id, name, highway
|
||||
HAVING COUNT(*) >= 3 -- Minimum 3 samples for reliability / 3 عينات على الأقل
|
||||
ORDER BY sample_count DESC
|
||||
`);
|
||||
let segmentsUpdated = 0;
|
||||
let totalPoints = 0;
|
||||
for (const row of matchedData) {
|
||||
const medianSpeed = parseFloat(row.median_speed);
|
||||
const sampleCount = parseInt(row.sample_count);
|
||||
const geojson = JSON.parse(row.geojson);
|
||||
totalPoints += sampleCount;
|
||||
let congestionFactor = 1.0;
|
||||
const majorRoadTypes = ['primary', 'secondary', 'trunk', 'motorway'];
|
||||
if (majorRoadTypes.includes(row.highway) && medianSpeed < 30) {
|
||||
congestionFactor = 30 / Math.max(medianSpeed, 1);
|
||||
}
|
||||
await this.roadStatRepo.upsert({
|
||||
segmentId: row.segment_id,
|
||||
averageSpeed: medianSpeed,
|
||||
congestionFactor,
|
||||
sampleCount,
|
||||
lastUpdated: new Date(),
|
||||
geometry: geojson,
|
||||
}, ['segmentId']);
|
||||
segmentsUpdated++;
|
||||
}
|
||||
return {
|
||||
segmentsUpdated,
|
||||
totalPointsProcessed: totalPoints,
|
||||
topAdjustments: matchedData.slice(0, 10).map(d => ({
|
||||
segmentId: d.segment_id,
|
||||
name: d.name || 'Unnamed Road',
|
||||
speed: Math.round(parseFloat(d.median_speed)),
|
||||
samples: parseInt(d.sample_count)
|
||||
}))
|
||||
};
|
||||
INSERT INTO road_segment_stats ("segmentId", "averageSpeed", "sampleCount", "lastUpdated", geometry, "congestionFactor")
|
||||
SELECT id, avg_speed, samples, NOW(), way4326,
|
||||
CASE WHEN highway IN ('primary','secondary','trunk','motorway') AND avg_speed < 30
|
||||
THEN 30/GREATEST(avg_speed,1) ELSE 1.0 END
|
||||
FROM stats
|
||||
ON CONFLICT ("segmentId") DO UPDATE SET
|
||||
"averageSpeed"=EXCLUDED."averageSpeed",
|
||||
"sampleCount"=EXCLUDED."sampleCount",
|
||||
"lastUpdated"=NOW(),
|
||||
"geometry"=EXCLUDED."geometry",
|
||||
"congestionFactor"=EXCLUDED."congestionFactor"
|
||||
RETURNING "segmentId";
|
||||
`;
|
||||
const result = await this.dataSource.query(query);
|
||||
return { segmentsUpdated: result.length };
|
||||
}
|
||||
async analyzeTimeProfiles(sinceHours = 720) {
|
||||
this.logger.log(`🕒 Starting temporal profiling (Phase 2) for last ${sinceHours}h...`);
|
||||
const query = `
|
||||
WITH grid_points AS (
|
||||
-- Bucket by 10m grid and Time (Hour + DOW)
|
||||
SELECT
|
||||
ST_SnapToGrid(ST_Transform(location::geometry, 3857), 10) AS loc,
|
||||
EXTRACT(HOUR FROM timestamp)::int as hr,
|
||||
EXTRACT(DOW FROM timestamp)::int as dow,
|
||||
AVG(speed) as speed
|
||||
FROM telemetry_logs
|
||||
WHERE timestamp >= NOW() - INTERVAL '${sinceHours} hours' AND speed > 2
|
||||
GROUP BY loc, hr, dow
|
||||
),
|
||||
matches AS (
|
||||
-- Map to OSM segments
|
||||
SELECT l.osm_id::text as sid, g.hr, g.dow, g.speed
|
||||
FROM grid_points g
|
||||
INNER JOIN planet_osm_line l ON
|
||||
l.highway IS NOT NULL AND
|
||||
l.way && ST_Expand(g.loc, 20) AND
|
||||
ST_DWithin(l.way, g.loc, 20)
|
||||
),
|
||||
temporal_stats AS (
|
||||
-- Aggregate by segment + time bucket
|
||||
SELECT sid, hr, dow,
|
||||
PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY speed) as avg_spd,
|
||||
COUNT(*) as smp
|
||||
FROM matches
|
||||
GROUP BY sid, hr, dow
|
||||
HAVING COUNT(*) >= 2 -- Require minimum samples per bucket
|
||||
)
|
||||
INSERT INTO road_speed_profiles ("segmentId", "hourOfDay", "dayOfWeek", "averageSpeed", "sampleCount", "lastUpdated")
|
||||
SELECT sid, hr, dow, avg_spd, smp, NOW()
|
||||
FROM temporal_stats
|
||||
ON CONFLICT ("segmentId", "hourOfDay", "dayOfWeek") DO UPDATE SET
|
||||
"averageSpeed" = EXCLUDED."averageSpeed",
|
||||
"sampleCount" = EXCLUDED."sampleCount",
|
||||
"lastUpdated" = NOW()
|
||||
RETURNING "segmentId";
|
||||
`;
|
||||
const result = await this.dataSource.query(query);
|
||||
this.logger.log(`📊 Time-aware profiling complete: ${result.length} buckets updated.`);
|
||||
return { bucketsUpdated: result.length };
|
||||
}
|
||||
async discoverNewRoads(sinceHours = 168) {
|
||||
this.logger.log(`🛣️ Starting road discovery for last ${sinceHours}h (${sinceHours / 24}d)...`);
|
||||
const candidates = await this.dataSource.query(`
|
||||
WITH off_road_points AS (
|
||||
SELECT
|
||||
t.id,
|
||||
t."driverId",
|
||||
t.speed,
|
||||
t.heading,
|
||||
t.location,
|
||||
t.timestamp
|
||||
FROM telemetry_logs t
|
||||
WHERE t.timestamp >= NOW() - INTERVAL '${sinceHours} hours'
|
||||
AND t.speed > 5 -- Moving, not parked / متحرك وليس متوقف
|
||||
AND NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM planet_osm_line l
|
||||
WHERE l.highway IS NOT NULL
|
||||
AND ST_DWithin(t.location::geography, ST_Transform(l.way, 4326)::geography, 15)
|
||||
)
|
||||
this.logger.log(`🛣️ Starting road discovery v2.5 for last ${sinceHours}h...`);
|
||||
const query = `
|
||||
WITH grid_points AS (
|
||||
-- Group by cell and driver to reduce volume for clustering and anti-join (v2.5.1 WoW Performance)
|
||||
SELECT "driverId", ST_SnapToGrid(ST_Transform(location::geometry, 3857), 10) as loc,
|
||||
AVG(speed) as speed,
|
||||
MIN(timestamp) as timestamp
|
||||
FROM telemetry_logs
|
||||
WHERE timestamp >= NOW() - INTERVAL '${sinceHours} hours' AND speed > 5
|
||||
GROUP BY loc, "driverId"
|
||||
),
|
||||
off_road AS (
|
||||
-- Spatial Anti-Join: Only points further than 35m from any existing highway
|
||||
SELECT g.* FROM grid_points g
|
||||
LEFT JOIN planet_osm_line l ON
|
||||
l.highway IS NOT NULL AND
|
||||
l.way && ST_Expand(g.loc, 35) AND
|
||||
ST_DWithin(l.way, g.loc, 35)
|
||||
WHERE l.osm_id IS NULL
|
||||
),
|
||||
-- Step 2: Cluster nearby off-road points using DBSCAN
|
||||
-- الخطوة 2: تجميع النقاط القريبة باستخدام DBSCAN
|
||||
clustered AS (
|
||||
SELECT
|
||||
*,
|
||||
ST_ClusterDBSCAN(location::geometry, eps := 0.0003, minpoints := 5)
|
||||
OVER () AS cluster_id
|
||||
FROM off_road_points
|
||||
-- Density-based clustering to find linear paths
|
||||
SELECT *, ST_ClusterDBSCAN(loc, eps := 30, minpoints := 3) OVER () as cid FROM off_road
|
||||
),
|
||||
cluster_stats AS (
|
||||
-- 1. Calculate cluster-wide quality metrics
|
||||
SELECT cid,
|
||||
COUNT(DISTINCT "driverId") as drv_count,
|
||||
COUNT(*) as pt_count,
|
||||
AVG(speed) as avg_spd
|
||||
FROM clustered
|
||||
WHERE cid IS NOT NULL
|
||||
GROUP BY cid
|
||||
HAVING COUNT(DISTINCT "driverId") >= 3 AND COUNT(*) >= 15
|
||||
),
|
||||
cluster_path AS (
|
||||
-- 2. Extract unique grid cells per cluster in chronological order
|
||||
SELECT cid, loc, MIN(timestamp) as ts
|
||||
FROM clustered
|
||||
WHERE cid IN (SELECT cid FROM cluster_stats)
|
||||
GROUP BY cid, loc
|
||||
),
|
||||
final_candidates AS (
|
||||
-- 3. Build geometry and calculate final scoring
|
||||
SELECT
|
||||
s.cid,
|
||||
ST_Transform(ST_MakeLine(p.loc ORDER BY p.ts), 4326) as geom,
|
||||
s.drv_count,
|
||||
s.pt_count,
|
||||
s.avg_spd,
|
||||
ST_Length(ST_Transform(ST_MakeLine(p.loc ORDER BY p.ts), 4326)::geography) as len
|
||||
FROM cluster_stats s
|
||||
JOIN cluster_path p ON s.cid = p.cid
|
||||
GROUP BY s.cid, s.drv_count, s.pt_count, s.avg_spd
|
||||
)
|
||||
-- Step 3: Aggregate clusters into candidate road lines
|
||||
-- الخطوة 3: تحويل التجمعات إلى خطوط طرق مرشحة
|
||||
SELECT
|
||||
cluster_id,
|
||||
COUNT(*) AS total_points,
|
||||
COUNT(DISTINCT "driverId") AS unique_drivers,
|
||||
AVG(speed) AS avg_speed,
|
||||
ST_AsGeoJSON(ST_MakeLine(location::geometry ORDER BY timestamp)) AS geojson_line,
|
||||
ST_Length(ST_MakeLine(location::geometry ORDER BY timestamp)::geography) AS length_m
|
||||
FROM clustered
|
||||
WHERE cluster_id IS NOT NULL
|
||||
GROUP BY cluster_id
|
||||
HAVING COUNT(DISTINCT "driverId") >= 2 -- At least 2 drivers / سائقان على الأقل
|
||||
AND COUNT(*) >= 10 -- At least 10 points / 10 نقاط على الأقل
|
||||
ORDER BY unique_drivers DESC, total_points DESC
|
||||
`);
|
||||
const savedCandidates = [];
|
||||
for (const c of candidates) {
|
||||
const geojson = JSON.parse(c.geojson_line);
|
||||
const confidence = this.calculateConfidence(parseInt(c.unique_drivers), parseInt(c.total_points), parseFloat(c.length_m));
|
||||
const candidate = this.candidateRoadRepo.create({
|
||||
geometry: geojson,
|
||||
uniqueDriverCount: parseInt(c.unique_drivers),
|
||||
totalPoints: parseInt(c.total_points),
|
||||
averageSpeed: parseFloat(c.avg_speed),
|
||||
lengthMeters: parseFloat(c.length_m),
|
||||
confidence,
|
||||
status: 'pending',
|
||||
});
|
||||
const saved = await this.candidateRoadRepo.save(candidate);
|
||||
savedCandidates.push({
|
||||
id: saved.id,
|
||||
uniqueDrivers: saved.uniqueDriverCount,
|
||||
totalPoints: saved.totalPoints,
|
||||
averageSpeed: Math.round(saved.averageSpeed * 10) / 10,
|
||||
lengthMeters: Math.round(saved.lengthMeters),
|
||||
confidence: Math.round(saved.confidence * 100) / 100,
|
||||
geojson,
|
||||
});
|
||||
}
|
||||
this.logger.log(`✅ Road discovery complete: ${savedCandidates.length} candidates found`);
|
||||
return {
|
||||
candidatesFound: savedCandidates.length,
|
||||
candidates: savedCandidates,
|
||||
};
|
||||
INSERT INTO candidate_roads (geometry, "uniqueDriverCount", "totalPoints", "averageSpeed", "lengthMeters", confidence, status)
|
||||
SELECT geom, drv_count, pt_count, avg_spd, len,
|
||||
-- SQL port of calculateConfidence logic
|
||||
ROUND((
|
||||
LEAST(drv_count::float / 5, 1.0) * 0.4 +
|
||||
LEAST(pt_count::float / 50, 1.0) * 0.3 +
|
||||
CASE WHEN len BETWEEN 50 AND 2000 THEN 0.3 ELSE 0.1 END
|
||||
)::numeric, 2) as conf,
|
||||
'pending'
|
||||
FROM final_candidates
|
||||
RETURNING id;
|
||||
`;
|
||||
const result = await this.dataSource.query(query);
|
||||
return { candidatesFound: result.length };
|
||||
}
|
||||
async getCongestionData(bounds) {
|
||||
return this.dataSource.query(`
|
||||
SELECT
|
||||
rs."segmentId" AS segment_id,
|
||||
rs."averageSpeed" AS avg_speed,
|
||||
rs."congestionFactor" AS congestion_factor,
|
||||
rs."sampleCount" AS sample_count,
|
||||
ST_AsGeoJSON(rs.geometry) AS geojson
|
||||
SELECT rs."segmentId", rs."averageSpeed", rs."congestionFactor", rs."sampleCount", ST_AsGeoJSON(rs.geometry) as geojson
|
||||
FROM road_segment_stats rs
|
||||
WHERE rs.geometry IS NOT NULL
|
||||
AND ST_Intersects(
|
||||
rs.geometry::geometry,
|
||||
ST_MakeEnvelope($1, $2, $3, $4, 4326)
|
||||
)
|
||||
AND rs."sampleCount" >= 3
|
||||
WHERE rs.geometry IS NOT NULL AND ST_Intersects(rs.geometry::geometry, ST_MakeEnvelope($1, $2, $3, $4, 4326))
|
||||
ORDER BY rs."congestionFactor" DESC
|
||||
`, [bounds.west, bounds.south, bounds.east, bounds.north]);
|
||||
}
|
||||
async getAnalysisSummary() {
|
||||
const [telemetryCount] = await this.dataSource.query(`SELECT COUNT(*) as count FROM telemetry_logs`);
|
||||
const [roadStatCount] = await this.dataSource.query(`SELECT COUNT(*) as count FROM road_segment_stats`);
|
||||
const [candidateCount] = await this.dataSource.query(`SELECT COUNT(*) as count,
|
||||
COUNT(*) FILTER (WHERE status = 'pending') as pending,
|
||||
COUNT(*) FILTER (WHERE status = 'approved') as approved,
|
||||
COUNT(*) FILTER (WHERE status = 'rejected') as rejected
|
||||
FROM candidate_roads`);
|
||||
const [recentActivity] = await this.dataSource.query(`
|
||||
SELECT
|
||||
COUNT(*) as points_last_24h,
|
||||
COUNT(DISTINCT "driverId") as active_drivers_24h
|
||||
FROM telemetry_logs
|
||||
WHERE timestamp >= NOW() - INTERVAL '24 hours'
|
||||
`);
|
||||
const [tCount] = await this.dataSource.query('SELECT COUNT(*) as count FROM telemetry_logs');
|
||||
const [rCount] = await this.dataSource.query('SELECT COUNT(*) as count FROM road_segment_stats');
|
||||
const [cCount] = await this.dataSource.query('SELECT COUNT(*) as count, COUNT(*) FILTER (WHERE status=\'pending\') as p FROM candidate_roads');
|
||||
const [clCount] = await this.dataSource.query('SELECT COUNT(*) as count FROM road_segment_stats WHERE "isClosed" = true');
|
||||
return {
|
||||
telemetry: {
|
||||
totalPoints: parseInt(telemetryCount?.count || '0'),
|
||||
last24h: parseInt(recentActivity?.points_last_24h || '0'),
|
||||
activeDrivers24h: parseInt(recentActivity?.active_drivers_24h || '0'),
|
||||
},
|
||||
roadSegments: {
|
||||
analyzed: parseInt(roadStatCount?.count || '0'),
|
||||
},
|
||||
candidateRoads: {
|
||||
total: parseInt(candidateCount?.count || '0'),
|
||||
pending: parseInt(candidateCount?.pending || '0'),
|
||||
approved: parseInt(candidateCount?.approved || '0'),
|
||||
rejected: parseInt(candidateCount?.rejected || '0'),
|
||||
},
|
||||
telemetry: { total: tCount.count },
|
||||
roads: { analyzed: rCount.count, closed: clCount.count },
|
||||
candidates: { total: cCount.count, pending: cCount.p }
|
||||
};
|
||||
}
|
||||
calculateConfidence(uniqueDrivers, totalPoints, lengthMeters) {
|
||||
const driverScore = Math.min(uniqueDrivers / 5, 1.0) * 0.4;
|
||||
const pointScore = Math.min(totalPoints / 50, 1.0) * 0.3;
|
||||
let lengthScore = 0;
|
||||
if (lengthMeters >= 50 && lengthMeters <= 2000) {
|
||||
lengthScore = 0.3;
|
||||
}
|
||||
else if (lengthMeters > 2000) {
|
||||
lengthScore = 0.2;
|
||||
}
|
||||
return Math.round((driverScore + pointScore + lengthScore) * 100) / 100;
|
||||
async getCandidates(status = 'pending', limit = 50) {
|
||||
return this.candidateRoadRepo.find({
|
||||
where: { status },
|
||||
order: { confidence: 'DESC' },
|
||||
take: limit
|
||||
});
|
||||
}
|
||||
async updateCandidateStatus(id, status) {
|
||||
await this.candidateRoadRepo.update(id, {
|
||||
status,
|
||||
reviewedAt: new Date()
|
||||
});
|
||||
return { success: true, id, status };
|
||||
}
|
||||
async discoverRoadClosures(sinceHours = 48) {
|
||||
this.logger.log(`🚧 Analyzing road closures for last ${sinceHours}h...`);
|
||||
await this.roadStatRepo.update({ isClosed: true }, { isClosed: false });
|
||||
const query = `
|
||||
WITH active_area AS (
|
||||
-- Bounding box of some recent activity to prove drivers are on the map
|
||||
SELECT ST_Expand(ST_Extent(location::geometry), 0.01) as bbox
|
||||
FROM telemetry_logs
|
||||
WHERE timestamp >= NOW() - INTERVAL '${sinceHours} hours'
|
||||
),
|
||||
possible_closures AS (
|
||||
SELECT rs."segmentId"
|
||||
FROM road_segment_stats rs
|
||||
WHERE rs."sampleCount" > 50
|
||||
AND rs.geometry && (SELECT bbox FROM active_area)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM telemetry_logs t
|
||||
WHERE t.timestamp >= NOW() - INTERVAL '${sinceHours} hours'
|
||||
AND ST_DWithin(rs.geometry::geometry, t.location::geometry, 35)
|
||||
)
|
||||
)
|
||||
UPDATE road_segment_stats
|
||||
SET "isClosed" = true
|
||||
WHERE "segmentId" IN (SELECT "segmentId" FROM possible_closures)
|
||||
RETURNING "segmentId";
|
||||
`;
|
||||
const result = await this.dataSource.query(query);
|
||||
this.logger.log(`✅ Road closure detection complete: ${result.length} roads flagged as closed.`);
|
||||
return { roadsClosed: result.length };
|
||||
}
|
||||
async getClosures() {
|
||||
return this.roadStatRepo.find({
|
||||
where: { isClosed: true },
|
||||
order: { sampleCount: 'DESC' }
|
||||
});
|
||||
}
|
||||
calculateConfidence(drivers, points, len) {
|
||||
const dScore = Math.min(drivers / 5, 1.0) * 0.4;
|
||||
const pScore = Math.min(points / 50, 1.0) * 0.3;
|
||||
const lScore = (len >= 50 && len <= 2000) ? 0.3 : 0.1;
|
||||
return Math.round((dScore + pScore + lScore) * 100) / 100;
|
||||
}
|
||||
};
|
||||
exports.TelemetryAnalyzerService = TelemetryAnalyzerService;
|
||||
__decorate([
|
||||
(0, schedule_1.Cron)('0 3 * * *'),
|
||||
(0, schedule_1.Cron)('0 4 */10 * *'),
|
||||
__metadata("design:type", Function),
|
||||
__metadata("design:paramtypes", []),
|
||||
__metadata("design:returntype", Promise)
|
||||
], TelemetryAnalyzerService.prototype, "handleNightlyIntelligence", null);
|
||||
__metadata("design:returntype", void 0)
|
||||
], TelemetryAnalyzerService.prototype, "handleDeepIntelligence", null);
|
||||
exports.TelemetryAnalyzerService = TelemetryAnalyzerService = TelemetryAnalyzerService_1 = __decorate([
|
||||
(0, common_1.Injectable)(),
|
||||
__param(0, (0, typeorm_1.InjectRepository)(telemetry_entity_1.TelemetryLog)),
|
||||
@@ -326,8 +377,10 @@ exports.TelemetryAnalyzerService = TelemetryAnalyzerService = TelemetryAnalyzerS
|
||||
__metadata("design:paramtypes", [typeorm_2.Repository,
|
||||
typeorm_2.Repository,
|
||||
typeorm_2.Repository,
|
||||
typeorm_2.DataSource,
|
||||
external_telemetry_service_1.ExternalTelemetryService,
|
||||
redis_service_1.RedisService,
|
||||
external_telemetry_service_1.ExternalTelemetryService])
|
||||
telegram_service_1.TelegramService,
|
||||
traffic_grid_service_1.TrafficGridService,
|
||||
typeorm_2.DataSource])
|
||||
], TelemetryAnalyzerService);
|
||||
//# sourceMappingURL=telemetry-analyzer.service.js.map
|
||||
Reference in New Issue
Block a user