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
33 changes: 33 additions & 0 deletions cmd/pggraph/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package main

import (
"fmt"
"os"

// Packages
pggraphcmd "github.com/mutablelogic/go-pg/pggraph/cmd"
servercmd "github.com/mutablelogic/go-server/pkg/cmd"
version "github.com/mutablelogic/go-server/pkg/version"
)

///////////////////////////////////////////////////////////////////////////////
// TYPES

type CLI struct {
pggraphcmd.ServerCommands // pggraph run
}

///////////////////////////////////////////////////////////////////////////////
// GLOBALS

const description = "PostgreSQL Graph is an application for running DAG graphs"

///////////////////////////////////////////////////////////////////////////////
// LIFECYCLE

func main() {
if err := servercmd.Main(CLI{}, description, version.Version()); err != nil {
fmt.Fprintln(os.Stderr, "Error:", err)
os.Exit(1)
}
}
88 changes: 88 additions & 0 deletions pggraph/cmd/server.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
package cmd

import (
"context"
"fmt"

// Packages
pg "github.com/mutablelogic/go-pg"
manager "github.com/mutablelogic/go-pg/pggraph/manager"
"github.com/mutablelogic/go-pg/pggraph/schema"
pgpkg "github.com/mutablelogic/go-pg/pkg/cmd"
server "github.com/mutablelogic/go-server"
cmd "github.com/mutablelogic/go-server/pkg/cmd"
"github.com/mutablelogic/go-server/pkg/types"
errgroup "golang.org/x/sync/errgroup"
)

///////////////////////////////////////////////////////////////////////////////
// TYPES

type ServerCommands struct {
Run RunServer `cmd:"" name:"run" help:"Run the server." group:"SERVER"`
}

type RunServer struct {
cmd.RunServer
pgpkg.PostgresFlags `embed:"" prefix:"pg."`
}

///////////////////////////////////////////////////////////////////////////////
// PUBLIC METHODS

func (runner *RunServer) Run(ctx server.Cmd) error {
// Connect to the database, if configured
conn, err := runner.PostgresFlags.Connect(ctx)
if err != nil {
return err
} else if conn == nil {
return fmt.Errorf("database connection is required")
}

// Report that the server is running
ctx.Logger().Info("running server", "name", ctx.Name(), "version", ctx.Version())

// Create the manager, run the server, and return any error
return runner.WithManager(ctx, conn, func(manager *manager.Manager) error {
// Create an error context - which will cancel any other goroutine on exit
errgroup, errctx := errgroup.WithContext(ctx.Context())

// Add a node to the manager
_, err := manager.RegisterNode(ctx.Context(), "helloworld", schema.NodeMeta{
Description: types.Ptr("Prints a greeting message"),
}, func(ctx context.Context, name string) (string, error) {
return "Hello, " + name, nil
})
if err != nil {
return err
}

// Run the manager
errgroup.Go(func() error {
return manager.Run(errctx, ctx.Logger())
})

// Run the server - if any co-routine in the error group returns an error, the server will be shutdown
errgroup.Go(func() error {
return runner.RunServer.Run(ctx.WithContext(errctx))
})

// Wait for the server and manager to exit, and return any error
return errgroup.Wait()
})
}

///////////////////////////////////////////////////////////////////////////////
// PRIVATE METHODS

func (runner *RunServer) WithManager(ctx server.Cmd, conn pg.PoolConn, fn func(*manager.Manager) error) error {
opts := []manager.Opt{
manager.WithMeter(ctx.Meter()),
manager.WithTracer(ctx.Tracer()),
}
if manager, err := manager.New(ctx.Context(), conn, ctx.Version(), opts...); err != nil {
return err
} else {
return fn(manager)
}
}
34 changes: 34 additions & 0 deletions pggraph/manager/graph.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package manager

import (
"context"

// Packages
otel "github.com/mutablelogic/go-client/pkg/otel"
pg "github.com/mutablelogic/go-pg"
schema "github.com/mutablelogic/go-pg/pggraph/schema"
types "github.com/mutablelogic/go-server/pkg/types"
attribute "go.opentelemetry.io/otel/attribute"
)

///////////////////////////////////////////////////////////////////////////////
// PUBLIC METHODS

// RegisterGraph creates a new graph, or updates an existing graph, and returns it.
func (manager *Manager) RegisterGraph(ctx context.Context, name string, meta schema.GraphMeta) (_ *schema.Graph, err error) {
ctx, endSpan := otel.StartSpan(manager.tracer, ctx, "RegisterGraph",
attribute.String("name", name),
attribute.String("meta", types.Stringify(meta)),
)
defer func() { endSpan(err) }()

var graph schema.Graph
if err := manager.queue.Tx(ctx, func(conn pg.Conn) error {
return conn.With("name", name, "version", manager.version).Insert(ctx, &graph, meta)
}); err != nil {
return nil, err
}

// Return success
return types.Ptr(graph), nil
}
102 changes: 102 additions & 0 deletions pggraph/manager/manager.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
package manager

import (
"context"
"fmt"
"strings"

// Packages
otel "github.com/mutablelogic/go-client/pkg/otel"
pg "github.com/mutablelogic/go-pg"
schema "github.com/mutablelogic/go-pg/pggraph/schema"
pgqueue "github.com/mutablelogic/go-pg/pgqueue/manager"
attribute "go.opentelemetry.io/otel/attribute"
)

///////////////////////////////////////////////////////////////////////////////
// TYPES

type Manager struct {
opt
queue *pgqueue.Manager
}

///////////////////////////////////////////////////////////////////////////////
// LIFECYCLE

func New(ctx context.Context, pool pg.PoolConn, version string, opts ...Opt) (*Manager, error) {
self := new(Manager)

// Check arguments
if pool == nil {
return nil, fmt.Errorf("pool is required")
}

// Set default values
if err := self.defaults(version); err != nil {
return nil, err
}

// Apply options
if err := self.apply(opts...); err != nil {
return nil, err
}

// Parse and register named queries so bind.Query(...) can resolve them.
queries, err := pg.NewQueries(strings.NewReader(schema.Queries))
if err != nil {
return nil, fmt.Errorf("parse queries.sql: %w", err)
} else {
pool = pool.WithQueries(queries).With(
"schema", self.schema,
).(pg.PoolConn)
}

// Create objects in the database schema. This is not done in a transaction
bootstrapCtx, endBootstrapSpan := otel.StartSpan(self.tracer, ctx, "bootstrap",
attribute.String("schema", self.schema),
)
if err := bootstrap(bootstrapCtx, pool, self.schema); err != nil {
endBootstrapSpan(err)
return nil, err
}

// Create the queue manager
queue, err := pgqueue.New(ctx, pool, pgqueue.WithSchema(self.schema), pgqueue.WithTracer(self.tracer), pgqueue.WithMeter(self.metrics), pgqueue.WithWorker(self.worker))
if err != nil {
endBootstrapSpan(err)
return nil, fmt.Errorf("create queue manager: %w", err)
} else {
self.queue = queue
}

// Return success
endBootstrapSpan(nil)
return self, nil
}

///////////////////////////////////////////////////////////////////////////////
// PRIVATE METHODS

func bootstrap(ctx context.Context, conn pg.Conn, schemaName string) error {
// Get all objects
objects, err := pg.NewQueries(strings.NewReader(schema.Objects))
if err != nil {
return fmt.Errorf("parse objects.sql: %w", err)
}

// Create the schema
if err := pg.SchemaCreate(ctx, conn, schemaName); err != nil {
return fmt.Errorf("create schema %q: %w", schemaName, err)
}

// Create all objects - not in a transaction
for _, key := range objects.Keys() {
if err := conn.Exec(ctx, objects.Query(key)); err != nil {
return fmt.Errorf("create object %q: %w", key, err)
}
}

// Return success
return nil
}
118 changes: 118 additions & 0 deletions pggraph/manager/node.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
package manager

import (
"context"
"fmt"
"reflect"
"runtime"
"strings"

// Packages
otel "github.com/mutablelogic/go-client/pkg/otel"
pg "github.com/mutablelogic/go-pg"
schema "github.com/mutablelogic/go-pg/pggraph/schema"
queueschema "github.com/mutablelogic/go-pg/pgqueue/schema"
types "github.com/mutablelogic/go-server/pkg/types"
attribute "go.opentelemetry.io/otel/attribute"
)

///////////////////////////////////////////////////////////////////////////////
// PUBLIC METHODS

// RegisterNode creates a new node, or updates an existing node, and returns it.
func (manager *Manager) RegisterNode(ctx context.Context, name string, meta schema.NodeMeta, fn any) (_ *schema.Node, err error) {
ctx, endSpan := otel.StartSpan(manager.tracer, ctx, "RegisterNode",
attribute.String("name", name),
attribute.String("meta", types.Stringify(meta)),
)
defer func() { endSpan(err) }()

// Determine the input and output edges from the function signature
_, in, out, err := functionEdges(fn)
if err != nil {
return nil, err
}

// Create the node in the database
insert := schema.NodeInsert{
Name: name,
In: out,
Out: in,
NodeMeta: meta,
}

var node schema.Node
if err := manager.queue.Tx(ctx, func(conn pg.Conn) error {
// Create the node
if err := conn.With("version", manager.version).Insert(ctx, &node, insert); err != nil {
return err
}

// TODO: Create the task queue for the node
_, err := manager.queue.RegisterQueue(ctx, name, queueschema.QueueMeta{}, manager.taskRunner(node, fn))
if err != nil {
return err
}

// Return success
return nil
}); err != nil {
return nil, err
}

// Return the node
return types.Ptr(node), nil
}

///////////////////////////////////////////////////////////////////////////////
// PRIVATE METHODS

func functionEdges(fn any) (string, []schema.Edge, schema.Edge, error) {
if fn == nil {
return "", nil, "", pg.ErrBadParameter.With("fn is missing")
}

value := reflect.ValueOf(fn)
if value.Kind() != reflect.Func {
return "", nil, "", pg.ErrBadParameter.With("fn is not a function")
}

fnType := value.Type()
if fnType.NumIn() < 2 {
return "", nil, "", pg.ErrBadParameter.With("fn must accept context.Context and at least one input")
}
if fnType.In(0) != reflect.TypeOf((*context.Context)(nil)).Elem() {
return "", nil, "", pg.ErrBadParameter.With("fn must accept context.Context as the first argument")
}
if fnType.NumOut() != 2 {
return "", nil, "", pg.ErrBadParameter.With("fn must return exactly one value and error")
}
if fnType.Out(1) != reflect.TypeOf((*error)(nil)).Elem() {
return "", nil, "", pg.ErrBadParameter.With("fn must return error as the last return value")
}

runtimeFn := runtime.FuncForPC(value.Pointer())
if runtimeFn == nil {
return "", nil, "", pg.ErrBadParameter.With("fn name is missing")
}
fnName := runtimeFn.Name()
if strings.TrimSpace(fnName) == "" {
return "", nil, "", pg.ErrBadParameter.With("fn name is missing")
}

in := make([]schema.Edge, 0, fnType.NumIn()-1)
for i := 1; i < fnType.NumIn(); i++ {
edge := schema.Edge(fnType.In(i).String())
if strings.TrimSpace(string(edge)) == "" {
return "", nil, "", fmt.Errorf("fn %q has an empty input type", fnName)
}
in = append(in, edge)
}

out := schema.Edge(fnType.Out(0).String())
if strings.TrimSpace(string(out)) == "" {
return "", nil, "", fmt.Errorf("fn %q has an empty return type", fnName)
}

return fnName, in, out, nil
}
Loading
Loading