Files

130 lines
3.2 KiB
Go

package dynamicProxyController
import (
"context"
"github.com/jmoiron/sqlx"
"github.com/michaelquigley/df/dl"
"github.com/openziti/sdk-golang/ziti"
"github.com/openziti/zrok/v2/controller/store"
"google.golang.org/grpc"
)
type Controller struct {
UnimplementedDynamicProxyControllerServer
str *store.Store
publisher *AmqpPublisher
zCfg *ziti.Config
zCtx ziti.Context
}
func NewController(cfg *Config, str *store.Store) (*Controller, error) {
publisher, err := NewAmqpPublisher(cfg.AmqpPublisher)
if err != nil {
return nil, err
}
zCfg, err := ziti.NewConfigFromFile(cfg.IdentityPath)
if err != nil {
return nil, err
}
zCtx, err := ziti.NewContext(zCfg)
if err != nil {
return nil, err
}
srv := grpc.NewServer()
ctrl := &Controller{
str: str,
publisher: publisher,
zCfg: zCfg,
zCtx: zCtx,
}
RegisterDynamicProxyControllerServer(srv, ctrl)
l, err := zCtx.Listen(cfg.ServiceName)
if err != nil {
return nil, err
}
go func() {
if err := srv.Serve(l); err != nil {
dl.Errorf("error serving dynamic proxy controller: %v", err)
return
}
}()
dl.Infof("started dynamic proxy controller server")
return ctrl, nil
}
func (c *Controller) BindFrontendMapping(frontendToken, name, shareToken string, trx *sqlx.Tx) error {
// create new frontend mapping
fm := &store.FrontendMapping{
FrontendToken: frontendToken,
Name: name,
ShareToken: shareToken,
}
fmId, err := c.str.CreateFrontendMapping(fm, trx)
if err != nil {
return err
}
// broadcast the mapping update via AMQP
mapping := Mapping{
Id: int64(fmId),
Operation: OperationBind,
Name: name,
ShareToken: shareToken,
}
return c.sendMappingUpdate(frontendToken, mapping)
}
func (c *Controller) UnbindFrontendMapping(frontendToken, name string, trx *sqlx.Tx) error {
if err := c.str.DeleteFrontendMappingsByFrontendTokenAndName(frontendToken, name, trx); err != nil {
return err
}
// broadcast the mapping update via AMQP
mapping := Mapping{
Operation: OperationUnbind,
Name: name,
}
return c.sendMappingUpdate(frontendToken, mapping)
}
func (c *Controller) FrontendMappings(_ context.Context, req *FrontendMappingsRequest) (*FrontendMappingsResponse, error) {
trx, err := c.str.Begin()
if err != nil {
return nil, err
}
defer trx.Rollback()
var mappings []*store.FrontendMapping
if req.GetName() == "" {
mappings, err = c.str.FindFrontendMappingsByFrontendTokenWithHigherId(req.GetFrontendToken(), req.GetId(), trx)
} else {
mappings, err = c.str.FindFrontendMappingsWithHigherId(req.GetFrontendToken(), req.GetName(), req.GetId(), trx)
}
if err != nil {
return nil, err
}
out := make([]*FrontendMapping, len(mappings))
for i, storeMapping := range mappings {
out[i] = &FrontendMapping{
Id: storeMapping.Id,
Name: storeMapping.Name,
ShareToken: storeMapping.ShareToken,
}
}
return &FrontendMappingsResponse{FrontendMappings: out}, nil
}
func (c *Controller) sendMappingUpdate(frontendToken string, m Mapping) error {
if err := c.publisher.Publish(context.Background(), frontendToken, m); err != nil {
return err
}
dl.Infof("sent mapping update '%+v' -> '%s'", m, frontendToken)
return nil
}