Skip to content
Merged
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
13 changes: 13 additions & 0 deletions connectors/kafka/kafkawrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -88,3 +88,16 @@ func (k *kafkaWrapConn) StreamLSN(ctx context.Context, r *connect.Request[adiomv
func (k *kafkaWrapConn) StreamUpdates(ctx context.Context, r *connect.Request[adiomv1.StreamUpdatesRequest], s *connect.ServerStream[adiomv1.StreamUpdatesResponse]) error {
return k.kafkaConn.StreamUpdates(ctx, r, s)
}

type teardownable interface {
Teardown()
}

func (k *kafkaWrapConn) Teardown() {
if t, ok := k.kafkaConn.(teardownable); ok {
t.Teardown()
}
if t, ok := k.kafkaWrapUnderlying.(teardownable); ok {
t.Teardown()
}
}
6 changes: 6 additions & 0 deletions connectors/mongo/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -969,6 +969,12 @@ func MongoClient(ctx context.Context, settings ConnectorSettings) (*mongo.Client

func (c *conn) Teardown() {
c.cancel()
c.buffersMutex.Lock()
for id, buf := range c.buffers {
buf.cleanup.Stop()
delete(c.buffers, id)
}
c.buffersMutex.Unlock()
_ = c.client.Disconnect(context.Background())
}

Expand Down
8 changes: 6 additions & 2 deletions connectors/testconn/connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -380,8 +380,12 @@ func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.Writ
}

func (c *conn) Teardown() {
_ = c.writeBootstrap.Close()
_ = c.writeUpdates.Close()
if c.writeBootstrap != nil {
_ = c.writeBootstrap.Close()
}
if c.writeUpdates != nil {
_ = c.writeUpdates.Close()
}
}

func NewConn(path string) adiomv1connect.ConnectorServiceHandler {
Expand Down
30 changes: 30 additions & 0 deletions internal/app/options/connectorflags.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,19 @@ import (
var ErrMissingConnector = errors.New("missing or unsupported connector")
var ErrHelp = errors.New("connector help used")

type teardownable interface {
Teardown()
}

func teardownConnector(cc ConfiguredConnector) {
if t, ok := cc.Local.(teardownable); ok {
t.Teardown()
}
if t, ok := cc.Remote.(teardownable); ok {
t.Teardown()
}
}

type AdditionalSettings struct {
BaseThreadCount int
}
Expand Down Expand Up @@ -90,17 +103,20 @@ func ConfigureConnectors(args []string, additionalSettings AdditionalSettings) (
}

if len(dstArgs) < 1 {
teardownConnector(src)
return src, dst, nil, fmt.Errorf("missing destination: %w", ErrMissingConnector)
}
var restArgs []string
for _, registeredConnector := range registeredConnectors {
if registeredConnector.IsConnector(dstArgs[0]) {
if registeredConnector.Create != nil {
if dst.Local, restArgs, err = registeredConnector.Create(dstArgs, additionalSettings); err != nil {
teardownConnector(src)
return src, dst, nil, err
}
} else {
if dst.Remote, restArgs, err = registeredConnector.CreateRemote(dstArgs, additionalSettings); err != nil {
teardownConnector(src)
return src, dst, nil, err
}
}
Expand All @@ -111,10 +127,12 @@ func ConfigureConnectors(args []string, additionalSettings AdditionalSettings) (
} else {
_, _, err = registeredConnector.CreateRemote([]string{dstArgs[0], "help"}, additionalSettings)
}
teardownConnector(src)
return src, dst, nil, err
}
}
if dst.Local == nil && dst.Remote == nil {
teardownConnector(src)
return src, dst, nil, fmt.Errorf("unsupported destination: %w", ErrMissingConnector)
}
return src, dst, restArgs, nil
Expand Down Expand Up @@ -552,13 +570,19 @@ func KafkaWrap(args []string, as AdditionalSettings) (adiomv1connect.ConnectorSe
}
registeredConnectors := GetRegisteredConnectors()
if len(restArgs) < 1 {
if t, ok := conn.(teardownable); ok {
t.Teardown()
}
return fmt.Errorf("missing connector for kafka-wrap")
}
for _, c := range registeredConnectors {
if c.IsConnector(restArgs[0]) {
if c.Create != nil {
underlying, newRestArgs, err := c.Create(restArgs, as)
if err != nil {
if t, ok := conn.(teardownable); ok {
t.Teardown()
}
return fmt.Errorf("err creating wrapped connector: %w", err)
}
restArgs = newRestArgs
Expand All @@ -567,12 +591,18 @@ func KafkaWrap(args []string, as AdditionalSettings) (adiomv1connect.ConnectorSe
} else if c.CreateRemote != nil {
underlying, newRestArgs, err := c.CreateRemote(restArgs, as)
if err != nil {
if t, ok := conn.(teardownable); ok {
t.Teardown()
}
return fmt.Errorf("err creating wrapped connector: %w", err)
}
restArgs = newRestArgs
conn = kafka.NewKafkaWrapConn(conn, underlying)
return nil
} else {
if t, ok := conn.(teardownable); ok {
t.Teardown()
}
return fmt.Errorf("unable to wrap connector")
}
}
Expand Down
Loading