629 lines
16 KiB
Go
629 lines
16 KiB
Go
package database
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
|
|
bin "gordenko.dev/dima/bin/little"
|
|
"gordenko.dev/dima/qb"
|
|
"gordenko.dev/dima/qb/bufreader"
|
|
"gordenko.dev/dima/qb/proto"
|
|
"gordenko.dev/dima/qb/storage"
|
|
"gordenko.dev/dima/qb/timeutil"
|
|
"gordenko.dev/dima/qb/transform"
|
|
"gordenko.dev/dima/qb/worker"
|
|
)
|
|
|
|
var (
|
|
ErrWrongResultCodeBug = errors.New("bug: wrong result code")
|
|
|
|
successMsg = []byte{
|
|
proto.RespSuccess,
|
|
}
|
|
)
|
|
|
|
func reply(conn io.Writer, errcode uint16) {
|
|
var answer []byte
|
|
if errcode == 0 {
|
|
answer = successMsg
|
|
} else {
|
|
answer = []byte{
|
|
proto.RespError,
|
|
0, 0,
|
|
}
|
|
bin.PutUint16(answer[1:], errcode)
|
|
}
|
|
_, err := conn.Write(answer)
|
|
if err != nil {
|
|
return
|
|
}
|
|
}
|
|
|
|
func (s *Database) handleTCPConn(conn net.Conn) {
|
|
defer conn.Close()
|
|
|
|
fmt.Println("client connected")
|
|
|
|
r := bufreader.New(conn, 128)
|
|
|
|
for {
|
|
err := s.processRequest(conn, r)
|
|
if err != nil {
|
|
if err != io.EOF {
|
|
s.logger.Println(err)
|
|
}
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (err error) {
|
|
messageType, err := r.ReadByte()
|
|
if err != nil {
|
|
if err != io.EOF {
|
|
return fmt.Errorf("read messageType: %s", err)
|
|
} else {
|
|
return err
|
|
}
|
|
}
|
|
|
|
switch messageType {
|
|
case proto.TypeGetMetric:
|
|
req, err := proto.ReadGetMetricReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadGetMetricReq: %s", err)
|
|
}
|
|
if err = s.GetMetric(conn, req); err != nil {
|
|
return fmt.Errorf("GetMetric: %s", err)
|
|
}
|
|
|
|
case proto.TypeAddMetric:
|
|
req, err := proto.ReadAddMetricReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadAddMetricReq: %s", err)
|
|
}
|
|
fmt.Println("ReadAddMetricReq:", req)
|
|
reply(conn, s.AddMetric(req))
|
|
|
|
case proto.TypeDeleteMetric:
|
|
req, err := proto.ReadDeleteMetricReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadDeleteMetricReq: %s", err)
|
|
}
|
|
reply(conn, s.DeleteMetric(req))
|
|
|
|
// case proto.TypeAppendMeasure:
|
|
// req, err := proto.ReadAppendMeasureReq(r)
|
|
// if err != nil {
|
|
// return fmt.Errorf("proto.ReadAppendMeasureReq: %s", err)
|
|
// }
|
|
// //fmt.Println("append measure", req.MetricID, conn.RemoteAddr().String())
|
|
// reply(conn, s.AppendMeasure(req))
|
|
|
|
case proto.TypeAppendMeasures:
|
|
req, err := proto.ReadAppendMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadAppendMeasuresReq: %s", err)
|
|
}
|
|
if err = s.AppendMeasures(conn, req); err != nil {
|
|
return fmt.Errorf("AppendMeasures: %s", err)
|
|
}
|
|
|
|
case proto.TypeListInstantMeasures:
|
|
req, err := proto.ReadListInstantMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListInstantMeasuresReq: %s", err)
|
|
}
|
|
if err = s.ListInstantMeasures(conn, req); err != nil {
|
|
return fmt.Errorf("ListInstantMeasures: %s", err)
|
|
}
|
|
|
|
case proto.TypeListCumulativeMeasures:
|
|
req, err := proto.ReadListCumulativeMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListCumulativeMeasuresReq: %s", err)
|
|
}
|
|
if err = s.ListCumulativeMeasures(conn, req); err != nil {
|
|
return fmt.Errorf("ListCumulativeMeasures: %s", err)
|
|
}
|
|
|
|
case proto.TypeListInstantPeriods:
|
|
req, err := proto.ReadListInstantPeriodsReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListInstantPeriodsReq: %s", err)
|
|
}
|
|
if err = s.ListInstantPeriods(conn, req); err != nil {
|
|
return fmt.Errorf("ListInstantPeriods: %s", err)
|
|
}
|
|
|
|
case proto.TypeListCumulativePeriods:
|
|
req, err := proto.ReadListCumulativePeriodsReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListCumulativePeriodsReq: %s", err)
|
|
}
|
|
if err = s.ListCumulativePeriods(conn, req); err != nil {
|
|
return fmt.Errorf("ListCumulativePeriods: %s", err)
|
|
}
|
|
|
|
case proto.TypeListCurrentValues:
|
|
req, err := proto.ReadListCurrentValuesReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ListCurrentValuesReq: %s", err)
|
|
}
|
|
if err = s.ListCurrentValues(conn, req); err != nil {
|
|
return fmt.Errorf("ListCurrentValues: %s", err)
|
|
}
|
|
|
|
case proto.TypeDeleteMeasures:
|
|
req, err := proto.ReadDeleteMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadDeleteMeasuresReq: %s", err)
|
|
}
|
|
reply(conn, s.DeleteMeasures(req))
|
|
|
|
case proto.TypeListAllInstantMeasures:
|
|
req, err := proto.ReadListAllInstantMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListAllInstantMeasuresReq: %s", err)
|
|
}
|
|
if err = s.ListAllInstantMeasures(conn, req); err != nil {
|
|
return fmt.Errorf("ListAllInstantMeasures: %s", err)
|
|
}
|
|
|
|
case proto.TypeListAllCumulativeMeasures:
|
|
req, err := proto.ReadListAllCumulativeMeasuresReq(r)
|
|
if err != nil {
|
|
return fmt.Errorf("proto.ReadListAllCumulativeMeasuresReq: %s", err)
|
|
}
|
|
if err = s.ListAllCumulativeMeasures(conn, req); err != nil {
|
|
return fmt.Errorf("ListAllCumulativeMeasures: %s", err)
|
|
}
|
|
|
|
default:
|
|
fmt.Printf("unknown messageType: %d\n", messageType)
|
|
return fmt.Errorf("unknown messageType: %d", messageType)
|
|
}
|
|
return
|
|
}
|
|
|
|
// API
|
|
|
|
func (s *Database) AddMetric(req proto.AddMetricReq) uint16 {
|
|
if req.MetricID == 0 {
|
|
return proto.ErrEmptyMetricID
|
|
}
|
|
if byte(req.FracDigits) > qb.MaxFracDigits {
|
|
return proto.ErrWrongFracDigits
|
|
}
|
|
switch req.MetricType {
|
|
case qb.Cumulative, qb.Instant:
|
|
// ok
|
|
default:
|
|
return proto.ErrWrongMetricType
|
|
}
|
|
|
|
resultCh := make(chan byte, 1)
|
|
|
|
s.workerInbox.Push(worker.AddMetricReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: req.MetricType,
|
|
FracDigits: req.FracDigits,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
resultCode := <-resultCh
|
|
|
|
switch resultCode {
|
|
case worker.Succeed:
|
|
//
|
|
case worker.MetricDuplicate:
|
|
return proto.ErrDuplicate
|
|
default:
|
|
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error {
|
|
resultCh := make(chan worker.GetMetricResult, 1)
|
|
|
|
s.workerInbox.Push(worker.GetMetricReq{
|
|
MetricID: req.MetricID,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
result := <-resultCh
|
|
|
|
switch result.ResultCode {
|
|
case worker.Succeed:
|
|
answer := []byte{
|
|
proto.RespValue,
|
|
0, 0, 0, 0, // metricID
|
|
byte(result.MetricType),
|
|
result.FracDigits,
|
|
}
|
|
bin.PutUint32(answer[1:], req.MetricID)
|
|
|
|
_, err := conn.Write(answer)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
case worker.NoMetric:
|
|
reply(conn, proto.ErrNoMetric)
|
|
default:
|
|
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 {
|
|
// resultCh := make(chan worker.DeleteMetricResult, 1)
|
|
|
|
// s.workerInbox.Push(worker.DeleteMetricReq{
|
|
// MetricID: req.MetricID,
|
|
// ResultCh: resultCh,
|
|
// })
|
|
|
|
// result := <-resultCh
|
|
|
|
// switch result.ResultCode {
|
|
// case worker.Succeed:
|
|
// // var (
|
|
// // //freePageNumbers []uint32
|
|
// // )
|
|
// // if result.RootPageNo > 0 {
|
|
// // var err error
|
|
// // freePageNumbers, err = s.atree.GetAllPages(result.RootPageNo)
|
|
// // if err != nil {
|
|
// // qb.Abort(qb.FailedAtreeRequest, err)
|
|
// // }
|
|
// // }
|
|
// // s.storage.Append(storage.DeletedMetric{
|
|
// // MetricID: req.MetricID,
|
|
// // //FreePageNumbers: freePageNumbers, FIX
|
|
// // })
|
|
// //<-waitCh
|
|
|
|
// case worker.NoMetric:
|
|
// return proto.ErrNoMetric
|
|
|
|
// default:
|
|
// qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
// }
|
|
return 0
|
|
}
|
|
|
|
type FilledPage struct {
|
|
Since uint32
|
|
RootPageNo uint32
|
|
PrevPageNo uint32
|
|
TimestampsChunks [][]byte
|
|
TimestampsSize uint16
|
|
ValuesChunks [][]byte
|
|
ValuesSize uint16
|
|
}
|
|
|
|
func (s *Database) AppendMeasures(conn io.Writer, req proto.AppendMeasuresReq) error {
|
|
resultCh := make(chan storage.MeasuresAppendResult, 1)
|
|
|
|
s.workerInbox.Push(worker.AppendMeasuresReq{
|
|
MetricID: req.MetricID,
|
|
Measures: req.Measures,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
result := <-resultCh
|
|
|
|
switch result.ResultCode {
|
|
case worker.Succeed, worker.NonMonotonicValue, worker.ExpiredMeasure:
|
|
var errorCode byte
|
|
switch result.ResultCode {
|
|
case worker.NonMonotonicValue:
|
|
errorCode = proto.ErrNonMonotonicValue
|
|
case worker.ExpiredMeasure:
|
|
errorCode = proto.ErrExpiredMeasure
|
|
}
|
|
answer := []byte{
|
|
proto.RespValue,
|
|
0, 0, // written count
|
|
errorCode,
|
|
}
|
|
bin.PutIntAsUint16(answer[1:], result.WrittenCount)
|
|
_, err := conn.Write(answer)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
case worker.NoMetric:
|
|
reply(conn, proto.ErrNoMetric)
|
|
default:
|
|
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 {
|
|
resultCh := make(chan byte, 1)
|
|
|
|
s.workerInbox.Push(worker.DeleteMeasuresReq{
|
|
MetricID: req.MetricID,
|
|
Since: req.Since,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
result := <-resultCh
|
|
|
|
_ = result
|
|
|
|
// switch result.ResultCode {
|
|
// case worker.NoMeasuresToDelete:
|
|
// // ok
|
|
|
|
// //case worker.DeleteFromAtreeNotNeeded:
|
|
// // регистрирую удаление в TransactionLog
|
|
// // s.storage.Append(storage.MeasuresDeleteRecord{
|
|
// // MetricID: req.MetricID,
|
|
// // })
|
|
// //<-waitCh
|
|
|
|
// //case worker.DeleteFromAtreeRequired:
|
|
// // собираю номера всех data и index страниц метрики (типа запись REDO лога).
|
|
// // pageNumbers, err := s.atree.GetAllPages(req.MetricID)
|
|
// // if err != nil {
|
|
// // qb.Abort(qb.FailedAtreeRequest, err)
|
|
// // }
|
|
// // // регистрирую удаление в TransactionLog
|
|
// // s.storage.Append(storage.DeletedMeasures{
|
|
// // MetricID: req.MetricID,
|
|
// // FreePageNumbers: pageNumbers,
|
|
// // })
|
|
// //<-waitCh
|
|
|
|
// case worker.NoMetric:
|
|
// return proto.ErrNoMetric
|
|
|
|
// default:
|
|
// qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
// }
|
|
return 0
|
|
}
|
|
|
|
// SELECT
|
|
|
|
func (s *Database) ListAllInstantMeasures(conn net.Conn, req proto.ListAllInstantMetricMeasuresReq) error {
|
|
responseWriter := transform.NewInstantMeasureWriter(conn, 0)
|
|
|
|
return s.fullScan(fullScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Instant,
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
func (s *Database) ListAllCumulativeMeasures(conn io.Writer, req proto.ListAllCumulativeMeasuresReq) error {
|
|
responseWriter := transform.NewCumulativeMeasureWriter(conn, 0)
|
|
|
|
return s.fullScan(fullScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Cumulative,
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
func (s *Database) ListInstantMeasures(conn net.Conn, req proto.ListInstantMeasuresReq) error {
|
|
if req.Since > req.Until {
|
|
reply(conn, proto.ErrInvalidRange)
|
|
return nil
|
|
}
|
|
|
|
responseWriter := transform.NewInstantMeasureWriter(conn, req.Since)
|
|
|
|
return s.rangeScan(rangeScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Instant,
|
|
Since: req.Since,
|
|
Until: req.Until - 1,
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
func (s *Database) ListCumulativeMeasures(conn net.Conn, req proto.ListCumulativeMeasuresReq) error {
|
|
if req.Since > req.Until {
|
|
reply(conn, proto.ErrInvalidRange)
|
|
return nil
|
|
}
|
|
|
|
responseWriter := transform.NewCumulativeMeasureWriter(conn, req.Since)
|
|
|
|
return s.rangeScan(rangeScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Cumulative,
|
|
Since: req.Since,
|
|
Until: req.Until - 1,
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriodsReq) error {
|
|
since, until := timeutil.TimeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay)
|
|
if since.After(until) {
|
|
reply(conn, proto.ErrInvalidRange)
|
|
return nil
|
|
}
|
|
responseWriter, err := transform.NewInstantPeriodsWriter(transform.InstantPeriodsWriterOptions{
|
|
Dst: conn,
|
|
GroupBy: req.GroupBy,
|
|
AggregateFuncs: req.AggregateFuncs,
|
|
FirstHourOfDay: req.FirstHourOfDay,
|
|
})
|
|
if err != nil {
|
|
reply(conn, proto.ErrUnexpected)
|
|
return nil
|
|
}
|
|
return s.rangeScan(rangeScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Instant,
|
|
Since: uint32(since.Unix()),
|
|
Until: uint32(until.Unix()),
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
func (s *Database) ListCumulativePeriods(conn net.Conn, req proto.ListCumulativePeriodsReq) error {
|
|
since, until := timeutil.TimeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay)
|
|
if since.After(until) {
|
|
reply(conn, proto.ErrInvalidRange)
|
|
return nil
|
|
}
|
|
responseWriter, err := transform.NewCumulativePeriodsWriter(transform.CumulativePeriodsWriterOptions{
|
|
Dst: conn,
|
|
GroupBy: req.GroupBy,
|
|
FirstHourOfDay: req.FirstHourOfDay,
|
|
})
|
|
if err != nil {
|
|
reply(conn, proto.ErrUnexpected)
|
|
return nil
|
|
}
|
|
return s.rangeScan(rangeScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: qb.Cumulative,
|
|
Since: uint32(since.Unix()),
|
|
Until: uint32(until.Unix()),
|
|
Conn: conn,
|
|
ResponseWriter: responseWriter,
|
|
})
|
|
}
|
|
|
|
type rangeScanReq struct {
|
|
MetricID uint32
|
|
MetricType qb.MetricType
|
|
Since uint32
|
|
Until uint32
|
|
Conn io.Writer
|
|
ResponseWriter qb.PeriodsWriter
|
|
}
|
|
|
|
func (s *Database) rangeScan(req rangeScanReq) error {
|
|
resultCh := make(chan worker.RangeScanResult, 1)
|
|
|
|
s.workerInbox.Push(worker.RangeScanReq{
|
|
MetricID: req.MetricID,
|
|
Since: req.Since,
|
|
Until: req.Until,
|
|
MetricType: req.MetricType,
|
|
ResponseWriter: req.ResponseWriter,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
result := <-resultCh
|
|
|
|
switch result.ResultCode {
|
|
case worker.QueryDone:
|
|
req.ResponseWriter.Close()
|
|
case worker.UntilFound:
|
|
err := s.ContinueRangeScan(ContinueRangeScanReq{
|
|
MetricType: req.MetricType,
|
|
FracDigits: result.FracDigits,
|
|
ResponseWriter: req.ResponseWriter,
|
|
LastPageNo: result.LastPageNo,
|
|
Since: req.Since,
|
|
})
|
|
//s.metricRUnlock(req.MetricID) fix release unlock
|
|
if err != nil {
|
|
reply(req.Conn, proto.ErrUnexpected)
|
|
} else {
|
|
req.ResponseWriter.Close()
|
|
}
|
|
case worker.UntilNotFound:
|
|
err := s.RangeScan(RangeScanReq{
|
|
MetricType: req.MetricType,
|
|
FracDigits: result.FracDigits,
|
|
ResponseWriter: req.ResponseWriter,
|
|
Since: req.Since,
|
|
Until: req.Until,
|
|
LastPageNo: result.LastPageNo,
|
|
IsDataPage: result.IsDataPage,
|
|
})
|
|
//s.metricRUnlock(req.MetricID) // fix release
|
|
if err != nil {
|
|
reply(req.Conn, proto.ErrUnexpected)
|
|
} else {
|
|
req.ResponseWriter.Close()
|
|
}
|
|
case worker.NoMetric:
|
|
reply(req.Conn, proto.ErrNoMetric)
|
|
case worker.WrongMetricType:
|
|
reply(req.Conn, proto.ErrWrongMetricType)
|
|
default:
|
|
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type fullScanReq struct {
|
|
MetricID uint32
|
|
MetricType qb.MetricType
|
|
Conn io.Writer
|
|
ResponseWriter qb.PeriodsWriter
|
|
}
|
|
|
|
func (s *Database) fullScan(req fullScanReq) error {
|
|
resultCh := make(chan worker.FullScanResult, 1)
|
|
|
|
s.workerInbox.Push(worker.FullScanReq{
|
|
MetricID: req.MetricID,
|
|
MetricType: req.MetricType,
|
|
ResponseWriter: req.ResponseWriter,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
result := <-resultCh
|
|
|
|
switch result.ResultCode {
|
|
case worker.QueryDone:
|
|
req.ResponseWriter.Close()
|
|
case worker.UntilFound:
|
|
err := s.ContinueFullScan(ContinueFullScanReq{
|
|
MetricType: req.MetricType,
|
|
FracDigits: result.FracDigits,
|
|
ResponseWriter: req.ResponseWriter,
|
|
LastPageNo: result.LastPageNo,
|
|
})
|
|
//s.worker.AddJobToQueue(req.MetricID) // FIX release rlock
|
|
if err != nil {
|
|
reply(req.Conn, proto.ErrUnexpected)
|
|
} else {
|
|
req.ResponseWriter.Close()
|
|
}
|
|
case worker.NoMetric:
|
|
reply(req.Conn, proto.ErrNoMetric)
|
|
case worker.WrongMetricType:
|
|
reply(req.Conn, proto.ErrWrongMetricType)
|
|
default:
|
|
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Database) ListCurrentValues(conn net.Conn, req proto.ListCurrentValuesReq) error {
|
|
responseWriter := transform.NewCurrentValueWriter(conn)
|
|
defer responseWriter.Close()
|
|
|
|
resultCh := make(chan struct{})
|
|
|
|
s.workerInbox.Push(worker.ListCurrentValuesReq{
|
|
MetricIDs: req.MetricIDs,
|
|
ResponseWriter: responseWriter,
|
|
ResultCh: resultCh,
|
|
})
|
|
|
|
<-resultCh
|
|
return nil
|
|
}
|