2019-08-08 14:47:04 +01:00
|
|
|
// Copyright (C) 2019 Storj Labs, Inc.
|
|
|
|
// See LICENSE for copying information.
|
|
|
|
|
|
|
|
package storagenodedb
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"database/sql"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
"github.com/zeebo/errs"
|
|
|
|
|
2019-12-27 11:48:47 +00:00
|
|
|
"storj.io/common/storj"
|
2021-04-23 10:52:40 +01:00
|
|
|
"storj.io/private/tagsql"
|
2019-08-08 14:47:04 +01:00
|
|
|
"storj.io/storj/storagenode/storageusage"
|
|
|
|
)
|
|
|
|
|
2019-09-18 17:17:28 +01:00
|
|
|
// StorageUsageDBName represents the database name.
|
|
|
|
const StorageUsageDBName = "storage_usage"
|
2019-08-21 15:32:25 +01:00
|
|
|
|
2020-07-16 15:18:02 +01:00
|
|
|
// storageUsageDB storage usage DB.
|
2019-09-18 17:17:28 +01:00
|
|
|
type storageUsageDB struct {
|
2019-11-13 16:49:22 +00:00
|
|
|
dbContainerImpl
|
2019-08-08 14:47:04 +01:00
|
|
|
}
|
|
|
|
|
2020-07-16 15:18:02 +01:00
|
|
|
// Store stores storage usage stamps to db replacing conflicting entries.
|
2019-09-18 17:17:28 +01:00
|
|
|
func (db *storageUsageDB) Store(ctx context.Context, stamps []storageusage.Stamp) (err error) {
|
2019-08-08 14:47:04 +01:00
|
|
|
defer mon.Task()(&ctx)(&err)
|
|
|
|
|
|
|
|
if len(stamps) == 0 {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2022-07-21 12:15:03 +01:00
|
|
|
query := `INSERT OR REPLACE INTO storage_usage(satellite_id, at_rest_total, interval_end_time, timestamp)
|
|
|
|
VALUES(?,?,?,?)`
|
2019-08-08 14:47:04 +01:00
|
|
|
|
2020-01-17 19:08:29 +00:00
|
|
|
return withTx(ctx, db.GetDB(), func(tx tagsql.Tx) error {
|
2019-08-08 14:47:04 +01:00
|
|
|
for _, stamp := range stamps {
|
2022-07-21 12:15:03 +01:00
|
|
|
_, err = tx.ExecContext(ctx, query, stamp.SatelliteID, stamp.AtRestTotal, stamp.IntervalEndTime.UTC(), stamp.IntervalStart.UTC())
|
2019-08-08 14:47:04 +01:00
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
}
|
|
|
|
|
|
|
|
// GetDaily returns daily storage usage stamps for particular satellite
|
2020-07-16 15:18:02 +01:00
|
|
|
// for provided time range.
|
2019-09-18 17:17:28 +01:00
|
|
|
func (db *storageUsageDB) GetDaily(ctx context.Context, satelliteID storj.NodeID, from, to time.Time) (_ []storageusage.Stamp, err error) {
|
2019-08-08 14:47:04 +01:00
|
|
|
defer mon.Task()(&ctx)(&err)
|
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
// hour_interval = current row interval_end_time - previous row interval_end_time
|
2022-07-21 12:15:03 +01:00
|
|
|
// Rows with 0-hour difference are assumed to be 24 hours.
|
2022-12-10 00:27:15 +00:00
|
|
|
query := `SELECT su1.satellite_id,
|
|
|
|
su1.at_rest_total,
|
2022-12-05 14:10:04 +00:00
|
|
|
COALESCE(
|
2022-12-10 00:27:15 +00:00
|
|
|
(
|
|
|
|
CAST(strftime('%s', su1.interval_end_time) AS NUMERIC)
|
|
|
|
-
|
|
|
|
CAST(strftime('%s', (
|
|
|
|
SELECT interval_end_time
|
|
|
|
FROM storage_usage
|
|
|
|
WHERE satellite_id = su1.satellite_id
|
|
|
|
AND timestamp < su1.timestamp
|
|
|
|
ORDER BY timestamp DESC
|
|
|
|
LIMIT 1
|
|
|
|
)) AS NUMERIC)
|
|
|
|
) / 3600,
|
|
|
|
24
|
2022-12-05 14:10:04 +00:00
|
|
|
) AS hour_interval,
|
2022-12-10 00:27:15 +00:00
|
|
|
su1.timestamp
|
|
|
|
FROM storage_usage su1
|
|
|
|
WHERE su1.satellite_id = ?
|
|
|
|
AND ? <= su1.timestamp AND su1.timestamp <= ?
|
|
|
|
ORDER BY su1.timestamp ASC`
|
2019-08-08 14:47:04 +01:00
|
|
|
|
2019-08-21 15:32:25 +01:00
|
|
|
rows, err := db.QueryContext(ctx, query, satelliteID, from.UTC(), to.UTC())
|
2019-08-08 14:47:04 +01:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2020-01-16 14:36:50 +00:00
|
|
|
defer func() { err = errs.Combine(err, rows.Close()) }()
|
2019-08-08 14:47:04 +01:00
|
|
|
|
|
|
|
var stamps []storageusage.Stamp
|
|
|
|
for rows.Next() {
|
|
|
|
var satellite storj.NodeID
|
2022-12-05 14:10:04 +00:00
|
|
|
var atRestTotal, intervalInHours float64
|
2022-07-21 12:15:03 +01:00
|
|
|
var timestamp time.Time
|
2019-08-08 14:47:04 +01:00
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
err = rows.Scan(&satellite, &atRestTotal, &intervalInHours, ×tamp)
|
2019-08-08 14:47:04 +01:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
stamps = append(stamps, storageusage.Stamp{
|
2022-12-05 14:10:04 +00:00
|
|
|
SatelliteID: satellite,
|
|
|
|
AtRestTotal: atRestTotal,
|
|
|
|
AtRestTotalBytes: atRestTotal / intervalInHours,
|
|
|
|
IntervalInHours: intervalInHours,
|
|
|
|
IntervalStart: timestamp,
|
2019-08-08 14:47:04 +01:00
|
|
|
})
|
|
|
|
}
|
|
|
|
|
2020-01-16 14:36:50 +00:00
|
|
|
return stamps, rows.Err()
|
2019-08-08 14:47:04 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
// GetDailyTotal returns daily storage usage stamps summed across all known satellites
|
2020-07-16 15:18:02 +01:00
|
|
|
// for provided time range.
|
2019-09-18 17:17:28 +01:00
|
|
|
func (db *storageUsageDB) GetDailyTotal(ctx context.Context, from, to time.Time) (_ []storageusage.Stamp, err error) {
|
2019-08-08 14:47:04 +01:00
|
|
|
defer mon.Task()(&ctx)(&err)
|
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
// hour_interval = current row interval_end_time - previous row interval_end_time
|
2022-07-21 12:15:03 +01:00
|
|
|
// Rows with 0-hour difference are assumed to be 24 hours.
|
2022-12-10 00:27:15 +00:00
|
|
|
query := `SELECT SUM(su3.at_rest_total), SUM(su3.hour_interval), su3.timestamp
|
2022-07-21 12:15:03 +01:00
|
|
|
FROM (
|
2022-12-10 00:27:15 +00:00
|
|
|
SELECT su1.at_rest_total,
|
|
|
|
COALESCE(
|
|
|
|
(
|
|
|
|
CAST(strftime('%s', su1.interval_end_time) AS NUMERIC)
|
|
|
|
-
|
|
|
|
CAST(strftime('%s', (
|
|
|
|
SELECT interval_end_time
|
|
|
|
FROM storage_usage su2
|
|
|
|
WHERE su2.satellite_id = su1.satellite_id
|
|
|
|
AND su2.timestamp < su1.timestamp
|
|
|
|
ORDER BY su2.timestamp DESC
|
|
|
|
LIMIT 1
|
|
|
|
)) AS NUMERIC)
|
|
|
|
) / 3600,
|
|
|
|
24
|
|
|
|
) AS hour_interval,
|
|
|
|
su1.timestamp
|
|
|
|
FROM storage_usage su1
|
|
|
|
WHERE ? <= su1.timestamp AND su1.timestamp <= ?
|
|
|
|
) as su3
|
|
|
|
GROUP BY su3.timestamp
|
|
|
|
ORDER BY su3.timestamp ASC`
|
2019-08-08 14:47:04 +01:00
|
|
|
|
2019-08-21 15:32:25 +01:00
|
|
|
rows, err := db.QueryContext(ctx, query, from.UTC(), to.UTC())
|
2019-08-08 14:47:04 +01:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
defer func() {
|
|
|
|
err = errs.Combine(err, rows.Close())
|
|
|
|
}()
|
|
|
|
|
|
|
|
var stamps []storageusage.Stamp
|
|
|
|
for rows.Next() {
|
2022-12-05 14:10:04 +00:00
|
|
|
var atRestTotal, intervalInHours float64
|
2022-07-21 12:15:03 +01:00
|
|
|
var timestamp time.Time
|
2019-08-08 14:47:04 +01:00
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
err = rows.Scan(&atRestTotal, &intervalInHours, ×tamp)
|
2019-08-08 14:47:04 +01:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
stamps = append(stamps, storageusage.Stamp{
|
2022-12-05 14:10:04 +00:00
|
|
|
AtRestTotal: atRestTotal,
|
|
|
|
AtRestTotalBytes: atRestTotal / intervalInHours,
|
|
|
|
IntervalInHours: intervalInHours,
|
|
|
|
IntervalStart: timestamp,
|
2019-08-08 14:47:04 +01:00
|
|
|
})
|
|
|
|
}
|
|
|
|
|
2020-01-16 14:36:50 +00:00
|
|
|
return stamps, rows.Err()
|
2019-08-08 14:47:04 +01:00
|
|
|
}
|
2019-08-21 15:32:25 +01:00
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
// Summary returns aggregated storage usage in Bytes*hour and average usage in bytes across all satellites.
|
|
|
|
func (db *storageUsageDB) Summary(ctx context.Context, from, to time.Time) (_, _ float64, err error) {
|
2019-09-04 15:13:43 +01:00
|
|
|
defer mon.Task()(&ctx, from, to)(&err)
|
2022-12-05 14:10:04 +00:00
|
|
|
var summary, averageUsageInBytes sql.NullFloat64
|
2019-09-04 15:13:43 +01:00
|
|
|
|
2022-12-10 00:27:15 +00:00
|
|
|
query := `SELECT SUM(su3.at_rest_total), AVG(su3.at_rest_total_bytes)
|
2022-10-10 12:35:58 +01:00
|
|
|
FROM (
|
|
|
|
SELECT
|
2022-12-05 14:10:04 +00:00
|
|
|
at_rest_total,
|
|
|
|
at_rest_total / (
|
|
|
|
COALESCE(
|
2022-12-10 00:27:15 +00:00
|
|
|
(
|
|
|
|
CAST(strftime('%s', su1.interval_end_time) AS NUMERIC)
|
|
|
|
-
|
|
|
|
CAST(strftime('%s', (
|
|
|
|
SELECT interval_end_time
|
|
|
|
FROM storage_usage su2
|
|
|
|
WHERE su2.satellite_id = su1.satellite_id
|
|
|
|
AND su2.timestamp < su1.timestamp
|
|
|
|
ORDER BY su2.timestamp DESC
|
|
|
|
LIMIT 1
|
|
|
|
)) AS NUMERIC)
|
|
|
|
) / 3600,
|
|
|
|
24
|
|
|
|
)
|
2022-12-05 14:10:04 +00:00
|
|
|
) AS at_rest_total_bytes
|
2022-12-10 00:27:15 +00:00
|
|
|
FROM storage_usage su1
|
2022-10-10 12:35:58 +01:00
|
|
|
WHERE ? <= timestamp AND timestamp <= ?
|
2022-12-10 00:27:15 +00:00
|
|
|
) as su3`
|
2019-09-04 15:13:43 +01:00
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
err = db.QueryRowContext(ctx, query, from.UTC(), to.UTC()).Scan(&summary, &averageUsageInBytes)
|
|
|
|
return summary.Float64, averageUsageInBytes.Float64, err
|
2019-09-04 15:13:43 +01:00
|
|
|
}
|
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
// SatelliteSummary returns aggregated storage usage in Bytes*hour and average usage in bytes for a particular satellite.
|
|
|
|
func (db *storageUsageDB) SatelliteSummary(ctx context.Context, satelliteID storj.NodeID, from, to time.Time) (_, _ float64, err error) {
|
2019-09-04 15:13:43 +01:00
|
|
|
defer mon.Task()(&ctx, satelliteID, from, to)(&err)
|
2022-12-05 14:10:04 +00:00
|
|
|
var summary, averageUsageInBytes sql.NullFloat64
|
2019-09-04 15:13:43 +01:00
|
|
|
|
2022-12-10 00:27:15 +00:00
|
|
|
query := `SELECT SUM(su3.at_rest_total), AVG(su3.at_rest_total_bytes)
|
2022-10-10 12:35:58 +01:00
|
|
|
FROM (
|
|
|
|
SELECT
|
2022-12-05 14:10:04 +00:00
|
|
|
at_rest_total,
|
|
|
|
at_rest_total / (
|
|
|
|
COALESCE(
|
2022-12-10 00:27:15 +00:00
|
|
|
(
|
|
|
|
CAST(strftime('%s', su1.interval_end_time) AS NUMERIC)
|
|
|
|
-
|
|
|
|
CAST(strftime('%s', (
|
|
|
|
SELECT interval_end_time
|
|
|
|
FROM storage_usage su2
|
|
|
|
WHERE su2.satellite_id = su1.satellite_id
|
|
|
|
AND su2.timestamp < su1.timestamp
|
|
|
|
ORDER BY su2.timestamp DESC
|
|
|
|
LIMIT 1
|
|
|
|
)) AS NUMERIC)
|
|
|
|
) / 3600,
|
|
|
|
24
|
|
|
|
)
|
2022-12-05 14:10:04 +00:00
|
|
|
) AS at_rest_total_bytes
|
2022-12-10 00:27:15 +00:00
|
|
|
FROM storage_usage su1
|
2022-10-10 12:35:58 +01:00
|
|
|
WHERE satellite_id = ?
|
|
|
|
AND ? <= timestamp AND timestamp <= ?
|
2022-12-10 00:27:15 +00:00
|
|
|
) as su3`
|
2019-09-04 15:13:43 +01:00
|
|
|
|
2022-12-05 14:10:04 +00:00
|
|
|
err = db.QueryRowContext(ctx, query, satelliteID, from.UTC(), to.UTC()).Scan(&summary, &averageUsageInBytes)
|
|
|
|
return summary.Float64, averageUsageInBytes.Float64, err
|
2019-09-04 15:13:43 +01:00
|
|
|
}
|