surrealpatch/db/db.go

558 lines
12 KiB
Go
Raw Normal View History

// Copyright © 2016 Abcum Ltd
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package db
import (
"context"
"fmt"
"io"
"time"
2016-10-27 08:33:14 +00:00
"net/http"
2016-07-21 21:49:47 +00:00
"runtime/debug"
"github.com/abcum/fibre"
"github.com/abcum/surreal/cnf"
2016-05-25 11:33:05 +00:00
"github.com/abcum/surreal/kvs"
"github.com/abcum/surreal/log"
"github.com/abcum/surreal/mem"
"github.com/abcum/surreal/sql"
2016-07-16 13:43:53 +00:00
"cloud.google.com/go/trace"
_ "github.com/abcum/surreal/kvs/rixxdb"
// _ "github.com/abcum/surreal/kvs/dendro"
)
2017-02-20 01:09:24 +00:00
var QueryNotExecuted = fmt.Errorf("Query not executed")
type Response struct {
2016-10-27 08:33:14 +00:00
Time string `codec:"time,omitempty"`
Status string `codec:"status,omitempty"`
Detail string `codec:"detail,omitempty"`
Result []interface{} `codec:"result,omitempty"`
}
var db *kvs.DS
// Setup sets up the connection with the data layer
func Setup(opts *cnf.Options) (err error) {
log.WithPrefix("db").Infof("Starting database")
2016-06-15 12:38:37 +00:00
db, err = kvs.New(opts)
2016-06-15 12:38:37 +00:00
return
}
// Exit shuts down the connection with the data layer
func Exit() {
2016-06-15 12:38:37 +00:00
log.WithPrefix("db").Infof("Gracefully shutting down database")
2016-05-25 11:33:05 +00:00
db.Close()
}
// Import loads database operations from a reader.
// This can be used to playback a database snapshot
// into an already running database.
func Import(r io.Reader) (err error) {
return db.Import(r)
}
// Export saves all database operations to a writer.
// This can be used to save a database snapshot
// to a secondary file or stream.
func Export(w io.Writer) (err error) {
return db.Export(w)
}
// Begin begins a new read / write transaction
// with the underlying database, and returns
// the transaction, or any error which occured.
func Begin(rw bool) (txn kvs.TX, err error) {
return db.Begin(rw)
}
// Execute parses a single sql query, or multiple
// sql queries, and executes them serially against
// the underlying data layer.
2016-09-14 21:32:52 +00:00
func Execute(ctx *fibre.Context, txt interface{}, vars map[string]interface{}) (out []*Response, err error) {
span := trace.FromContext(ctx.Context()).NewChild("db.Execute")
nctx := trace.NewContext(ctx.Context(), span)
ctx = ctx.WithContext(nctx)
defer span.Finish()
// If no preset variables have been defined
// then ensure that the variables is
// instantiated for future use.
if vars == nil {
vars = make(map[string]interface{})
}
// Parse the received SQL batch query strings
// into SQL ASTs, using any immutable preset
// variables if set.
2016-09-14 21:32:52 +00:00
ast, err := sql.Parse(ctx, txt, vars)
2016-02-27 12:05:35 +00:00
if err != nil {
return
}
// Ensure that the current authentication data
// is made available as a runtime variable to
// the query layer.
vars["auth"] = ctx.Get("auth").(*cnf.Auth).Data
return Process(ctx, ast, vars)
}
// Process takes a parsed set of sql queries and
// executes them serially against the underlying
// data layer.
func Process(ctx *fibre.Context, ast *sql.Query, vars map[string]interface{}) (out []*Response, err error) {
span := trace.FromContext(ctx.Context()).NewChild("db.Process")
nctx := trace.NewContext(ctx.Context(), span)
ctx = ctx.WithContext(nctx)
defer span.Finish()
// Create 2 channels, one for force quitting
// the query processor, and the other for
// receiving and buffering any query results.
quit := make(chan bool)
2016-10-27 08:33:14 +00:00
recv := make(chan *Response)
// Ensure that the force quit channel is auto
// closed when the end of the request has been
// reached, and we are not an http connection.
2016-10-27 08:33:14 +00:00
defer close(quit)
// If the current connection is a normal http
// connection then force quit any running
// queries if a socket close event occurs.
if _, ok := ctx.Response().Writer().(http.CloseNotifier); ok {
2016-10-27 08:33:14 +00:00
exit := ctx.Response().CloseNotify()
done := make(chan struct{})
2016-10-27 08:33:14 +00:00
defer close(done)
go func() {
select {
case <-done:
case <-exit:
quit <- true
}
}()
2016-10-27 08:33:14 +00:00
}
// Create a new query executor with the query
// details, and the current runtime variables
// and execute the queries within.
exec := newExec(ast, ctx, vars)
2017-02-14 20:02:18 +00:00
defer exec.done()
2017-02-14 20:02:18 +00:00
go exec.execute(ctx.Context(), quit, recv)
// Wait for all of the processed queries to
// return results, buffer the output, and
// return the output when finished.
2016-10-27 08:33:14 +00:00
for res := range recv {
out = append(out, res)
2016-02-27 12:05:35 +00:00
}
return
}
func (e *executor) execute(ctx context.Context, quit <-chan bool, send chan<- *Response) {
2016-10-27 08:33:14 +00:00
var err error
2017-02-09 20:45:45 +00:00
var now time.Time
2016-10-27 08:33:14 +00:00
var rsp *Response
var buf []*Response
var res []interface{}
// Ensure that the query responses channel is
// closed when the full query has been processed
// and dealt with.
defer close(send)
// If we are making use of a global transaction
// which is not committed at the end of the
// query set, then cancel the transaction.
defer func() {
if e.txn != nil {
e.txn.Cancel()
}
2016-10-27 08:33:14 +00:00
}()
2017-02-09 20:44:08 +00:00
// If we have panicked during query execution
2016-10-27 08:33:14 +00:00
// then ensure that we recover from the error
// and print the error to the log.
defer func() {
if err := recover(); err != nil {
stk := string(debug.Stack())
log.WithPrefix("db").WithFields(map[string]interface{}{
"ctx": e.web, "stack": stk,
}).Errorln(err)
}
}()
2016-07-04 10:37:29 +00:00
2016-10-27 08:33:14 +00:00
// Loop over the defined query statements and
// process them, while listening for the quit
// channel to see if the client has gone away.
for _, stm := range e.ast.Statements {
2016-10-18 12:49:46 +00:00
2016-10-27 08:33:14 +00:00
select {
2016-10-18 12:49:46 +00:00
2016-10-27 08:33:14 +00:00
case <-quit:
return
2016-10-18 12:49:46 +00:00
default:
trc := trace.FromContext(ctx).NewChild(fmt.Sprint(stm))
ctx := trace.NewContext(ctx, trc)
2016-10-27 08:33:14 +00:00
// If we are not inside a global transaction
// then reset the error to nil so that the
// next statement is not ignored.
if e.txn == nil {
2017-02-09 20:45:45 +00:00
err, now = nil, time.Now()
2016-10-27 08:33:14 +00:00
}
// When in debugging mode, log every sql
// query, along with the query execution
// speed, so we can analyse slow queries.
log.WithPrefix("sql").WithFields(map[string]interface{}{
2017-02-28 00:19:21 +00:00
"id": e.web.Get("id"),
"kind": e.web.Get("auth").(*cnf.Auth).Kind,
"auth": e.web.Get("auth").(*cnf.Auth).Data,
}).Debugln(stm)
2016-10-27 08:33:14 +00:00
// Check to see if the current statement is
// a TRANSACTION statement, and if it is
// then deal with it and move on to the next.
switch stm.(type) {
case *sql.BeginStatement:
err = e.begin(true)
trc.Finish()
2016-10-27 08:33:14 +00:00
continue
case *sql.CancelStatement:
err, buf = e.cancel(buf, err, send)
trc.Finish()
2016-10-27 08:33:14 +00:00
continue
case *sql.CommitStatement:
err, buf = e.commit(buf, err, send)
trc.Finish()
2016-10-27 08:33:14 +00:00
continue
}
// If an error has occured and we are inside
// a global transaction, then ignore all
// subsequent statements in the transaction.
2016-10-27 08:33:14 +00:00
if err == nil {
res, err = e.operate(ctx, stm)
2016-10-18 12:49:46 +00:00
} else {
2017-02-20 01:09:24 +00:00
res, err = []interface{}{}, QueryNotExecuted
2016-10-18 12:49:46 +00:00
}
2016-10-27 08:33:14 +00:00
rsp = &Response{
Time: time.Since(now).String(),
Status: status(err),
Detail: detail(err),
Result: append([]interface{}{}, res...),
}
// If we are not inside a global transaction
// then we can output the statement response
// immediately to the channel.
if e.txn == nil {
2016-10-27 08:33:14 +00:00
send <- rsp
}
// If we are inside a global transaction we
// must buffer the responses for output at
// the end of the transaction.
if e.txn != nil {
2016-10-29 11:29:20 +00:00
switch stm.(type) {
case *sql.ReturnStatement:
buf = clear(buf, rsp)
default:
buf = append(buf, rsp)
}
2016-10-27 08:33:14 +00:00
}
trc.Finish()
2016-10-18 12:49:46 +00:00
}
2016-10-27 08:33:14 +00:00
}
}
func (e *executor) operate(ctx context.Context, ast sql.Statement) (res []interface{}, err error) {
2016-10-27 08:33:14 +00:00
var loc bool
var trw bool
2016-10-27 08:33:14 +00:00
// If we are not inside a global transaction
// then grab a new transaction, ensuring that
// it is closed at the end.
if e.txn == nil {
2016-10-27 08:33:14 +00:00
loc = true
switch ast.(type) {
case *sql.InfoStatement:
trw = false
err = e.begin(trw)
2016-10-27 08:33:14 +00:00
default:
trw = true
err = e.begin(trw)
2016-10-27 08:33:14 +00:00
}
if err != nil {
return
}
defer e.txn.Cancel()
2016-10-27 08:33:14 +00:00
}
// Mark the beginning of this statement so we
// can monitor the running time, and ensure
// it runs no longer than specified.
if stm, ok := ast.(sql.KillableStatement); ok {
stm.Begin()
}
2016-10-27 08:33:14 +00:00
// Execute the defined statement, receiving the
// result set, and any errors which occured
// while processing the query.
switch stm := ast.(type) {
case *sql.InfoStatement:
res, err = e.executeInfoStatement(ctx, stm)
case *sql.LetStatement:
res, err = e.executeLetStatement(ctx, stm)
2016-10-29 11:29:20 +00:00
case *sql.ReturnStatement:
res, err = e.executeReturnStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.SelectStatement:
res, err = e.executeSelectStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.CreateStatement:
res, err = e.executeCreateStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.UpdateStatement:
res, err = e.executeUpdateStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.DeleteStatement:
res, err = e.executeDeleteStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.RelateStatement:
res, err = e.executeRelateStatement(ctx, stm)
case *sql.DefineNamespaceStatement:
res, err = e.executeDefineNamespaceStatement(ctx, stm)
case *sql.RemoveNamespaceStatement:
res, err = e.executeRemoveNamespaceStatement(ctx, stm)
case *sql.DefineDatabaseStatement:
res, err = e.executeDefineDatabaseStatement(ctx, stm)
case *sql.RemoveDatabaseStatement:
res, err = e.executeRemoveDatabaseStatement(ctx, stm)
case *sql.DefineLoginStatement:
res, err = e.executeDefineLoginStatement(ctx, stm)
case *sql.RemoveLoginStatement:
res, err = e.executeRemoveLoginStatement(ctx, stm)
case *sql.DefineTokenStatement:
res, err = e.executeDefineTokenStatement(ctx, stm)
case *sql.RemoveTokenStatement:
res, err = e.executeRemoveTokenStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.DefineScopeStatement:
res, err = e.executeDefineScopeStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.RemoveScopeStatement:
res, err = e.executeRemoveScopeStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.DefineTableStatement:
res, err = e.executeDefineTableStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.RemoveTableStatement:
res, err = e.executeRemoveTableStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.DefineFieldStatement:
res, err = e.executeDefineFieldStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.RemoveFieldStatement:
res, err = e.executeRemoveFieldStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.DefineIndexStatement:
res, err = e.executeDefineIndexStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
case *sql.RemoveIndexStatement:
res, err = e.executeRemoveIndexStatement(ctx, stm)
2016-10-27 08:33:14 +00:00
}
// If this is a local transaction for only the
// current statement, then commit or cancel
// depending on the result error.
if loc && !e.txn.Closed() {
if !trw || err != nil {
e.txn.Cancel()
e.txn = nil
2016-10-27 08:33:14 +00:00
} else {
e.txn.Commit()
e.txn = nil
2016-10-27 08:33:14 +00:00
}
}
// The statement has successfully cancelled
// or committed, so stop all the transaction
// timeout timers if any were set.
if stm, ok := ast.(sql.KillableStatement); ok {
stm.Cease()
}
2016-10-27 08:33:14 +00:00
return
}
func status(e error) (s string) {
switch e.(type) {
default:
return "OK"
case *kvs.DBError:
return "ERR_DB"
case *kvs.KVError:
return "ERR_KV"
case error:
return "ERR"
}
}
func detail(e error) (s string) {
switch err := e.(type) {
default:
return
case error:
return err.Error()
}
}
func clear(buf []*Response, rsp *Response) []*Response {
for i := len(buf) - 1; i >= 0; i-- {
buf[len(buf)-1] = nil
buf = buf[:len(buf)-1]
}
return append(buf, rsp)
}
func (e *executor) begin(rw bool) (err error) {
if e.txn == nil {
e.txn, err = db.Begin(rw)
e.mem = mem.New(e.txn)
}
return
}
func (e *executor) cancel(buf []*Response, err error, chn chan<- *Response) (error, []*Response) {
defer func() {
e.txn = nil
e.mem = nil
}()
if e.txn == nil {
return nil, buf
}
e.txn.Cancel()
for _, v := range buf {
v.Status = "ERR"
chn <- v
}
for i := len(buf) - 1; i >= 0; i-- {
buf[len(buf)-1] = nil
buf = buf[:len(buf)-1]
}
return nil, buf
}
func (e *executor) commit(buf []*Response, err error, chn chan<- *Response) (error, []*Response) {
defer func() {
e.txn = nil
e.mem = nil
}()
if e.txn == nil {
return nil, buf
}
if err != nil {
e.txn.Cancel()
} else {
e.txn.Commit()
}
for _, v := range buf {
if err != nil {
v.Status = "ERR"
}
chn <- v
}
for i := len(buf) - 1; i >= 0; i-- {
buf[len(buf)-1] = nil
buf = buf[:len(buf)-1]
}
return nil, buf
}