package integration
import (
"context"
"database/sql"
"net"
"sync"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"k8s.io/klog/v2"
"github.com/google/trillian"
"github.com/google/trillian/extension"
"github.com/google/trillian/log"
"github.com/google/trillian/quota"
"github.com/google/trillian/server"
"github.com/google/trillian/server/admin"
"github.com/google/trillian/server/interceptor"
"github.com/google/trillian/storage/mysql"
"github.com/google/trillian/storage/testdb"
"github.com/google/trillian/util/clock"
_ "github.com/go-sql-driver/mysql"
)
var (
sequencerWindow = time.Duration(0)
batchSize = 50
SequencerInterval = 500 * time.Millisecond
timeSource = clock.System
)
type LogEnv struct {
registry extension.Registry
pendingTasks *sync.WaitGroup
grpcServer *grpc.Server
adminServer *admin.Server
logServer *server.TrillianLogRPCServer
LogOperation log.Operation
Sequencer *log.OperationManager
sequencerCancel context.CancelFunc
ClientConn *grpc.ClientConn
Address string
Log trillian.TrillianLogClient
Admin trillian.TrillianAdminClient
DB *sql.DB
dbDone func(context.Context)
}
func NewLogEnv(ctx context.Context, numSequencers int, _ string) (*LogEnv, error) {
return NewLogEnvWithGRPCOptions(ctx, numSequencers, nil, nil)
}
func NewLogEnvWithGRPCOptions(ctx context.Context, numSequencers int, serverOpts []grpc.ServerOption, clientOpts []grpc.DialOption) (*LogEnv, error) {
db, done, err := testdb.NewTrillianDB(ctx, testdb.DriverMySQL)
if err != nil {
return nil, err
}
registry := extension.Registry{
AdminStorage: mysql.NewAdminStorage(db),
LogStorage: mysql.NewLogStorage(db, nil),
QuotaManager: quota.Noop(),
}
ret, err := NewLogEnvWithRegistryAndGRPCOptions(ctx, numSequencers, registry, serverOpts, clientOpts)
if err != nil {
if err := db.Close(); err != nil {
return nil, err
}
return nil, err
}
ret.DB = db
ret.dbDone = done
return ret, nil
}
func NewLogEnvWithRegistry(ctx context.Context, numSequencers int, registry extension.Registry) (*LogEnv, error) {
return NewLogEnvWithRegistryAndGRPCOptions(ctx, numSequencers, registry, nil, nil)
}
func NewLogEnvWithRegistryAndGRPCOptions(ctx context.Context, numSequencers int, registry extension.Registry, serverOpts []grpc.ServerOption, clientOpts []grpc.DialOption) (*LogEnv, error) {
serverOpts = append(serverOpts, grpc.UnaryInterceptor(interceptor.ErrorWrapper))
grpcServer := grpc.NewServer(serverOpts...)
adminServer := admin.New(registry, nil)
trillian.RegisterTrillianAdminServer(grpcServer, adminServer)
logServer := server.NewTrillianLogRPCServer(registry, timeSource)
trillian.RegisterTrillianLogServer(grpcServer, logServer)
sequencerManager := log.NewSequencerManager(registry, sequencerWindow)
var wg sync.WaitGroup
var sequencerTask *log.OperationManager
ctx, cancel := context.WithCancel(ctx)
info := log.OperationInfo{
Registry: registry,
BatchSize: batchSize,
NumWorkers: numSequencers,
RunInterval: SequencerInterval,
TimeSource: timeSource,
}
sequencerTask = log.NewOperationManager(info, sequencerManager)
wg.Add(1)
go func(wg *sync.WaitGroup, om *log.OperationManager) {
defer wg.Done()
om.OperationLoop(ctx)
}(&wg, sequencerTask)
addr, lis, err := listen()
if err != nil {
cancel()
return nil, err
}
wg.Add(1)
go func(wg *sync.WaitGroup, grpcServer *grpc.Server, lis net.Listener) {
defer wg.Done()
if err := grpcServer.Serve(lis); err != nil {
klog.Errorf("gRPC server stopped: %v", err)
klog.Flush()
}
}(&wg, grpcServer, lis)
if clientOpts == nil {
clientOpts = []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
}
cc, err := grpc.Dial(addr, clientOpts...)
if err != nil {
cancel()
return nil, err
}
return &LogEnv{
registry: registry,
pendingTasks: &wg,
grpcServer: grpcServer,
adminServer: adminServer,
logServer: logServer,
Address: addr,
ClientConn: cc,
Log: trillian.NewTrillianLogClient(cc),
Admin: trillian.NewTrillianAdminClient(cc),
LogOperation: sequencerManager,
Sequencer: sequencerTask,
sequencerCancel: cancel,
}, nil
}
func (env *LogEnv) Close() {
if env.sequencerCancel != nil {
env.sequencerCancel()
}
if err := env.ClientConn.Close(); err != nil {
klog.Errorf("ClientConn.Close(): %v", err)
}
env.grpcServer.GracefulStop()
env.pendingTasks.Wait()
if env.dbDone != nil {
env.dbDone(context.TODO())
}
}