Skip to content
Draft
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
14 changes: 13 additions & 1 deletion app/db.server.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,10 @@
import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js'
import {
drizzle,
type PostgresJsDatabase,
type PostgresJsQueryResultHKT,
} from 'drizzle-orm/postgres-js'
import { type ExtractTablesWithRelations } from 'drizzle-orm'
import { type PgTransaction } from 'drizzle-orm/pg-core'
import postgres, { type Sql } from 'postgres'
import invariant from 'tiny-invariant'
import * as schema from './db/schema'
Expand Down Expand Up @@ -57,3 +63,9 @@ function parsePoolSize(value: string | undefined): number {
}

export { drizzleClient, pg }

export type DatabaseTransaction = PgTransaction<
PostgresJsQueryResultHKT,
typeof schema,
ExtractTablesWithRelations<typeof schema>
>
88 changes: 49 additions & 39 deletions app/db/models/device.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ import {
type Sensor,
} from '~/db/schema'
import type * as schema from '~/db/schema/index'
import { drizzleClient } from '~/db.server'
import { drizzleClient, type DatabaseTransaction } from '~/db.server'
import BaseNewDeviceEmail, {
messages as BaseNewDeviceMessages,
} from '~/emails/base-new-device'
Expand Down Expand Up @@ -155,56 +155,66 @@ export type DeviceForSingleMeasurementWrite = Awaited<
ReturnType<typeof getDeviceForSingleMeasurementWrite>
>

export function getDeviceForMeasurementWrite({ id }: Pick<Device, 'id'>) {
return drizzleClient.query.device.findFirst({
where: (device, { eq }) => eq(device.id, id),
columns: {
id: true,
archivedAt: true,
useAuth: true,
apiKey: true,
},
with: {
sensors: {
columns: {
id: true,
title: true,
sensorType: true,
},
},
},
})
export async function getDeviceForMeasurementWrite(
{ id }: Pick<Device, 'id'>,
tx: DatabaseTransaction,
) {
const currentDevice = await lockDeviceForMeasurementWrite(id, tx)
if (!currentDevice) return undefined

const sensors = await tx
.select({
id: sensor.id,
title: sensor.title,
sensorType: sensor.sensorType,
})
.from(sensor)
.where(eq(sensor.deviceId, id))
.orderBy(sensor.id)
.for('update')

return { ...currentDevice, sensors }
}

export async function getDeviceForSingleMeasurementWrite({
id,
sensorId,
}: Pick<Device, 'id'> & { sensorId: string }) {
const [row] = await drizzleClient
export async function getDeviceForSingleMeasurementWrite(
{ id, sensorId }: Pick<Device, 'id'> & { sensorId: string },
tx: DatabaseTransaction,
) {
const currentDevice = await lockDeviceForMeasurementWrite(id, tx)
if (!currentDevice) return undefined

const [currentSensor] = await tx
.select({
id: sensor.id,
})
.from(sensor)
.where(and(eq(sensor.deviceId, currentDevice.id), eq(sensor.id, sensorId)))
.limit(1)
.for('update')

return {
...currentDevice,
sensors: currentSensor ? [currentSensor] : [],
}
}

async function lockDeviceForMeasurementWrite(
id: Device['id'],
tx: DatabaseTransaction,
) {
const [currentDevice] = await tx
.select({
id: device.id,
archivedAt: device.archivedAt,
useAuth: device.useAuth,
apiKey: device.apiKey,
sensorId: sensor.id,
})
.from(device)
.leftJoin(
sensor,
and(eq(sensor.deviceId, device.id), eq(sensor.id, sensorId)),
)
.where(eq(device.id, id))
.limit(1)
.for('share')

if (!row) return undefined

return {
id: row.id,
archivedAt: row.archivedAt,
useAuth: row.useAuth,
apiKey: row.apiKey,
sensors: row.sensorId ? [{ id: row.sensorId }] : [],
}
return currentDevice
}

export function getUserDevice({ id, userId }: Pick<Device, 'id' | 'userId'>) {
Expand Down
78 changes: 27 additions & 51 deletions app/db/models/measurement.server.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { and, desc, eq, gt, gte, inArray, lt, lte, sql } from 'drizzle-orm'
import { ArchivedDeviceError } from './device.server'
import {
type LastMeasurement,
location,
Expand All @@ -9,9 +8,8 @@
measurements1hourView,
measurements1monthView,
measurements1yearView,
device,
} from '~/db/schema'
import { drizzleClient } from '~/db.server'
import { drizzleClient, type DatabaseTransaction } from '~/db.server'
import {
type MinimalDevice,
type MeasurementWithLocation,
Expand Down Expand Up @@ -178,12 +176,12 @@
}

export async function saveMeasurements(
tx: DatabaseTransaction,
minimalDevice: MinimalDevice,
measurements: MeasurementWithLocation[],
timing?: MeasurementTiming | null,
): Promise<void> {
if (!device) throw new Error('No device given!')
if (!Array.isArray(measurements)) throw new Error('Array expected')

Check failure on line 184 in app/db/models/measurement.server.ts

View workflow job for this annotation

GitHub Actions / ⚡ Test

tests/routes/api.boxes.$deviceId.locations.spec.ts > openSenseMap API Routes: /api/boxes/:deviceId/locations

Error: Array expected ❯ saveMeasurements app/db/models/measurement.server.ts:184:42 ❯ tests/routes/api.boxes.$deviceId.locations.spec.ts:110:9

const sensorIds = new Set(minimalDevice.sensors.map((s: any) => s.id))
const lastMeasurements: Record<string, NonNullable<LastMeasurement>> = {}
Expand Down Expand Up @@ -234,57 +232,35 @@
locationUpdateCount: deviceLocationUpdates.length,
})

await drizzleClient.transaction(async (tx) => {
const [currentDevice] = await tx
.select({
id: device.id,
archivedAt: device.archivedAt,
})
.from(device)
.where(eq(device.id, minimalDevice.id))
.limit(1)
timing?.mark('transactionDeviceLookup')

if (!currentDevice) {
const error = new Error('Device not found')
error.name = 'NotFoundError'
throw error
}

if (currentDevice.archivedAt) {
throw new ArchivedDeviceError(currentDevice.id)
}

const locations =
deviceLocationUpdates.length > 0
? await findOrCreateLocations(deviceLocationUpdates)
: []
timing?.mark('findOrCreateLocations', {
locationCount: locations.length,
})

if (deviceLocationUpdates.length > 0) {
await addLocationUpdates(
deviceLocationUpdates,
minimalDevice.id,
locations,
)
}
timing?.mark('addLocationUpdates')
const locations =
deviceLocationUpdates.length > 0
? await findOrCreateLocations(deviceLocationUpdates, tx)
: []
timing?.mark('findOrCreateLocations', {
locationCount: locations.length,
})

await insertMeasurementsWithLocation(
measurements,
locations,
if (deviceLocationUpdates.length > 0) {
await addLocationUpdates(
deviceLocationUpdates,
minimalDevice.id,
locations,
tx,
{ shouldReturn: false },
timing,
)
timing?.mark('insertMeasurements')
await updateLastMeasurements(lastMeasurements, tx, timing)
timing?.mark('updateLastMeasurements')
})
timing?.mark('transaction')
}
timing?.mark('addLocationUpdates')

await insertMeasurementsWithLocation(
measurements,
locations,
minimalDevice.id,
tx,
{ shouldReturn: false },
timing,
)
timing?.mark('insertMeasurements')
await updateLastMeasurements(lastMeasurements, tx, timing)
timing?.mark('updateLastMeasurements')
}

export async function insertMeasurements(measurements: any[]): Promise<void> {
Expand Down
Loading
Loading