Skip to content

Commit 08de0b4

Browse files
committed
refactor: merge MountOperation and MountRequest
1 parent 3e0f3a7 commit 08de0b4

26 files changed

Lines changed: 274 additions & 334 deletions

cmd/mount-proxy-client/main.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,8 @@ import (
77

88
flag "github.com/spf13/pflag"
99

10-
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/proxy"
1110
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/proxy/client"
11+
mounterutils "github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
1212
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/utils"
1313
)
1414

@@ -21,7 +21,7 @@ func main() {
2121
utils.AddGoFlags(flag.CommandLine)
2222
flag.Parse()
2323

24-
var req proxy.MountRequest
24+
var req mounterutils.MountRequest
2525
err := json.NewDecoder(os.Stdin).Decode(&req)
2626
if err != nil {
2727
fmt.Println(err.Error())

pkg/mounter/adaptor_mounter.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package mounter
33
import (
44
"context"
55

6+
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
67
mountutils "k8s.io/mount-utils"
78
)
89

@@ -18,9 +19,9 @@ func NewAdaptorMounter(inner mountutils.Interface) Mounter {
1819
}
1920
}
2021

21-
func (m *AdaptorMounter) ExtendedMount(_ context.Context, op *MountOperation) error {
22-
if op == nil {
22+
func (m *AdaptorMounter) ExtendedMount(_ context.Context, req *utils.MountRequest) error {
23+
if req == nil {
2324
return nil
2425
}
25-
return m.Mount(op.Source, op.Target, op.FsType, op.Options)
26+
return m.Mount(req.Source, req.Target, req.Fstype, req.Options)
2627
}

pkg/mounter/cmd_mounter.go

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import (
77
"os/exec"
88
"time"
99

10-
mounterutils "github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
10+
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
1111
"k8s.io/mount-utils"
1212
)
1313

@@ -29,12 +29,12 @@ func NewOssCmdMounter(execPath, volumeID string, inner mount.Interface) Mounter
2929
}
3030
}
3131

32-
func (m *OssCmdMounter) ExtendedMount(ctx context.Context, op *MountOperation) error {
33-
if op == nil {
32+
func (m *OssCmdMounter) ExtendedMount(ctx context.Context, req *utils.MountRequest) error {
33+
if req == nil {
3434
return nil
3535
}
3636

37-
cmd := exec.CommandContext(ctx, m.execPath, getArgs(op)...)
37+
cmd := exec.CommandContext(ctx, m.execPath, getArgs(req)...)
3838
cmd.Stdout = os.Stdout
3939
cmd.Stderr = os.Stderr
4040

@@ -44,17 +44,17 @@ func (m *OssCmdMounter) ExtendedMount(ctx context.Context, op *MountOperation) e
4444
return nil
4545
}
4646

47-
func getArgs(op *MountOperation) []string {
48-
if op == nil {
47+
func getArgs(req *utils.MountRequest) []string {
48+
if req == nil {
4949
return nil
5050
}
51-
switch op.FsType {
52-
case mounterutils.OssFsType:
53-
return mount.MakeMountArgs(op.Source, op.Target, "", op.Options)
54-
case mounterutils.OssFs2Type:
55-
args := []string{"mount", op.Target}
56-
args = append(args, op.Args...)
57-
for _, o := range op.Options {
51+
switch req.Fstype {
52+
case utils.OssFsType:
53+
return mount.MakeMountArgs(req.Source, req.Target, "", req.Options)
54+
case utils.OssFs2Type:
55+
args := []string{"mount", req.Target}
56+
args = append(args, req.Args...)
57+
for _, o := range req.Options {
5858
args = append(args, fmt.Sprintf("--%s", o))
5959
}
6060
return args

pkg/mounter/connector_mounter.go

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,23 +3,24 @@ package mounter
33
import (
44
"context"
55

6+
mounterutils "github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
67
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/utils"
78
"k8s.io/klog/v2"
8-
mountutils "k8s.io/mount-utils"
9+
"k8s.io/mount-utils"
910
)
1011

1112
type ConnectorMounter struct {
1213
mounterPath string
13-
mountutils.Interface
14+
mount.Interface
1415
}
1516

1617
var _ Mounter = &ConnectorMounter{}
1718

18-
func (m *ConnectorMounter) ExtendedMount(_ context.Context, op *MountOperation) error {
19-
if op == nil {
19+
func (m *ConnectorMounter) ExtendedMount(_ context.Context, req *mounterutils.MountRequest) error {
20+
if req == nil {
2021
return nil
2122
}
22-
args := mountutils.MakeMountArgs(op.Source, op.Target, op.FsType, op.Options)
23+
args := mount.MakeMountArgs(req.Source, req.Target, req.Fstype, req.Options)
2324
mntCmd := []string{"systemd-run", "--scope", "--"}
2425
if m.mounterPath == "" {
2526
mntCmd = append(mntCmd, "mount")
@@ -35,15 +36,15 @@ func (m *ConnectorMounter) ExtendedMount(_ context.Context, op *MountOperation)
3536
}
3637

3738
func (m *ConnectorMounter) Mount(source string, target string, fstype string, options []string) error {
38-
return m.ExtendedMount(context.Background(), &MountOperation{
39+
return m.ExtendedMount(context.Background(), &mounterutils.MountRequest{
3940
Source: source,
4041
Target: target,
41-
FsType: fstype,
42+
Fstype: fstype,
4243
Options: options,
4344
})
4445
}
4546

46-
func NewConnectorMounter(inner mountutils.Interface, mounterPath string) Mounter {
47+
func NewConnectorMounter(inner mount.Interface, mounterPath string) Mounter {
4748
return &ConnectorMounter{
4849
mounterPath: mounterPath,
4950
Interface: inner,

pkg/mounter/interceptors/alinas_secret.go

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ func NewAlinasSecretInterceptor(regionID string) (*AlinasSecretInterceptor, erro
6868
}, nil
6969
}
7070

71-
func (i *AlinasSecretInterceptor) Intercept(ctx context.Context, op *mounter.MountOperation, handler mounter.MountHandler) (err error) {
71+
func (i *AlinasSecretInterceptor) Intercept(ctx context.Context, op *mounterutils.MountRequest, handler mounter.MountHandler) (err error) {
7272
if op == nil || op.AuthConfig == nil {
7373
return handler(ctx, op)
7474
}
@@ -102,7 +102,7 @@ func (i *AlinasSecretInterceptor) Intercept(ctx context.Context, op *mounter.Mou
102102
return
103103
}
104104

105-
func (i *AlinasSecretInterceptor) startTokenRefreshLoop(op *mounter.MountOperation) {
105+
func (i *AlinasSecretInterceptor) startTokenRefreshLoop(op *mounterutils.MountRequest) {
106106
if op == nil || op.Target == "" || op.AuthConfig == nil {
107107
return
108108
}
@@ -147,7 +147,7 @@ func (i *AlinasSecretInterceptor) isMountPoint(target string) bool {
147147
return !notMnt
148148
}
149149

150-
func (i *AlinasSecretInterceptor) refreshAndSaveRRSAToken(op *mounter.MountOperation) (credFilePath string, err error) {
150+
func (i *AlinasSecretInterceptor) refreshAndSaveRRSAToken(op *mounterutils.MountRequest) (credFilePath string, err error) {
151151
if op == nil || op.AuthConfig == nil || op.AuthConfig.AuthType != rrsaAuthType {
152152
return
153153
}
@@ -172,7 +172,7 @@ func (i *AlinasSecretInterceptor) refreshAndSaveRRSAToken(op *mounter.MountOpera
172172
return save(token)
173173
}
174174

175-
func (i *AlinasSecretInterceptor) refreshRRSAToken(op *mounter.MountOperation) (token ramRoleToken, err error) {
175+
func (i *AlinasSecretInterceptor) refreshRRSAToken(op *mounterutils.MountRequest) (token ramRoleToken, err error) {
176176
if op == nil {
177177
return
178178
}
@@ -226,7 +226,7 @@ func (i *AlinasSecretInterceptor) shouldRefreshRRSAToken(roleName string) bool {
226226
return now.After(token.refreshAt.Add(5*time.Minute)) || token.expiresAt.Before(now.Add(10*time.Minute))
227227
}
228228

229-
func saveCredentials(op *mounter.MountOperation) (credFilePath string, err error) {
229+
func saveCredentials(op *mounterutils.MountRequest) (credFilePath string, err error) {
230230
if op == nil || op.AuthConfig == nil {
231231
return
232232
}
@@ -248,7 +248,9 @@ func saveCredentials(op *mounter.MountOperation) (credFilePath string, err error
248248

249249
credFileContent := makeCredFileContent(op.AuthConfig)
250250
if _, err = tmpCredFile.Write(credFileContent); err != nil {
251-
tmpCredFile.Close()
251+
if err := tmpCredFile.Close(); err != nil {
252+
klog.ErrorS(err, "Failed to close temporary alinas credential file", "path", tmpFilePath)
253+
}
252254
return
253255
}
254256

pkg/mounter/interceptors/alinas_secret_test.go

Lines changed: 13 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@ import (
1515
"github.com/alibabacloud-go/tea/tea"
1616
"github.com/golang/mock/gomock"
1717
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/cloud"
18-
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter"
1918
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
2019
"github.com/stretchr/testify/assert"
2120
"github.com/stretchr/testify/require"
@@ -24,10 +23,10 @@ import (
2423
)
2524

2625
var (
27-
successMountHandler = func(context.Context, *mounter.MountOperation) error {
26+
successMountHandler = func(context.Context, *utils.MountRequest) error {
2827
return nil
2928
}
30-
failureMountHandler = func(context.Context, *mounter.MountOperation) error {
29+
failureMountHandler = func(context.Context, *utils.MountRequest) error {
3130
return fmt.Errorf("failed")
3231
}
3332
)
@@ -186,7 +185,7 @@ func TestSaveCredentials(t *testing.T) {
186185

187186
tests := []struct {
188187
name string
189-
op *mounter.MountOperation
188+
op *utils.MountRequest
190189
expectErr bool
191190
}{
192191
{
@@ -196,12 +195,12 @@ func TestSaveCredentials(t *testing.T) {
196195
},
197196
{
198197
name: "nil auth config",
199-
op: &mounter.MountOperation{},
198+
op: &utils.MountRequest{},
200199
expectErr: false,
201200
},
202201
{
203202
name: "valid credentials",
204-
op: &mounter.MountOperation{
203+
op: &utils.MountRequest{
205204
VolumeID: "volume-id",
206205
AuthConfig: &utils.AuthConfig{
207206
AccessKey: "test-ak",
@@ -285,7 +284,7 @@ func TestRefreshRRSAToken(t *testing.T) {
285284

286285
tests := []struct {
287286
name string
288-
op *mounter.MountOperation
287+
op *utils.MountRequest
289288
setupMock func(*cloud.MockSTSInterface)
290289
expectErr bool
291290
expectToken bool
@@ -297,7 +296,7 @@ func TestRefreshRRSAToken(t *testing.T) {
297296
},
298297
{
299298
name: "successful token refresh",
300-
op: &mounter.MountOperation{
299+
op: &utils.MountRequest{
301300
VolumeID: "volume-id",
302301
Target: "/mnt/target",
303302
AuthConfig: &utils.AuthConfig{
@@ -325,7 +324,7 @@ func TestRefreshRRSAToken(t *testing.T) {
325324
},
326325
{
327326
name: "sts client returns error",
328-
op: &mounter.MountOperation{
327+
op: &utils.MountRequest{
329328
VolumeID: "volume-id",
330329
Target: "/mnt/target",
331330
AuthConfig: &utils.AuthConfig{
@@ -394,7 +393,7 @@ func TestRefreshAndSaveRRSAToken(t *testing.T) {
394393

395394
tests := []struct {
396395
name string
397-
op *mounter.MountOperation
396+
op *utils.MountRequest
398397
existingToken *ramRoleToken
399398
setupMock func(*cloud.MockSTSInterface)
400399
expectRefresh bool
@@ -407,14 +406,14 @@ func TestRefreshAndSaveRRSAToken(t *testing.T) {
407406
},
408407
{
409408
name: "nil auth config",
410-
op: &mounter.MountOperation{
409+
op: &utils.MountRequest{
411410
VolumeID: "volume-id",
412411
},
413412
expectRefresh: false,
414413
},
415414
{
416415
name: "wrong auth type",
417-
op: &mounter.MountOperation{
416+
op: &utils.MountRequest{
418417
VolumeID: "volume-id",
419418
AuthConfig: &utils.AuthConfig{
420419
AuthType: "other",
@@ -424,7 +423,7 @@ func TestRefreshAndSaveRRSAToken(t *testing.T) {
424423
},
425424
{
426425
name: "token is fresh, should not refresh",
427-
op: &mounter.MountOperation{
426+
op: &utils.MountRequest{
428427
VolumeID: "volume-id",
429428
AuthConfig: &utils.AuthConfig{
430429
AuthType: rrsaAuthType,
@@ -442,7 +441,7 @@ func TestRefreshAndSaveRRSAToken(t *testing.T) {
442441
},
443442
{
444443
name: "token needs refresh, should create credential file",
445-
op: &mounter.MountOperation{
444+
op: &utils.MountRequest{
446445
VolumeID: "volume-id",
447446
Target: "/mnt/target",
448447
AuthConfig: &utils.AuthConfig{

pkg/mounter/interceptors/ossfs_monitor.go

Lines changed: 12 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import (
66

77
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter"
88
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/proxy/server"
9+
"github.com/kubernetes-sigs/alibaba-cloud-csi-driver/pkg/mounter/utils"
910
"k8s.io/klog/v2"
1011
"k8s.io/mount-utils"
1112
)
@@ -17,36 +18,36 @@ var (
1718
monitorManager = server.NewMountMonitorManager()
1819
)
1920

20-
func OssfsMonitorInterceptor(ctx context.Context, op *mounter.MountOperation, handler mounter.MountHandler) error {
21-
if op == nil || op.MetricsPath == "" {
22-
return handler(ctx, op)
21+
func OssfsMonitorInterceptor(ctx context.Context, req *utils.MountRequest, handler mounter.MountHandler) error {
22+
if req == nil || req.MetricsPath == "" {
23+
return handler(ctx, req)
2324
}
2425

2526
// Get or create monitor for this target
26-
monitor, found := monitorManager.GetMountMonitor(op.Target, op.MetricsPath, raw, true)
27+
monitor, found := monitorManager.GetMountMonitor(req.Target, req.MetricsPath, raw, true)
2728
if monitor == nil {
28-
klog.ErrorS(errors.New("failed to get mount monitor"), "stop monitoring mountpoint status", "mountpoint", op.Target)
29-
return handler(ctx, op)
29+
klog.ErrorS(errors.New("failed to get mount monitor"), "stop monitoring mountpoint status", "mountpoint", req.Target)
30+
return handler(ctx, req)
3031
}
3132
if found {
3233
monitor.IncreaseMountRetryCount()
3334
}
3435

35-
err := handler(ctx, op)
36+
err := handler(ctx, req)
3637

3738
if err != nil {
3839
// This method should only be called when err != nil.
3940
// Invoking it with a nil error will trigger a warning log.
4041
monitor.HandleMountFailureOrExit(err)
4142
}
4243

43-
if op.MountResult == nil {
44+
if req.MountResult == nil {
4445
return err
4546
}
4647

47-
res, ok := op.MountResult.(server.OssfsMountResult)
48+
res, ok := req.MountResult.(server.OssfsMountResult)
4849
if !ok {
49-
klog.ErrorS(errors.New("failed to assert ossfs mount result type"), "skipping monitoring of mountpoint", "mountpoint", op.Target)
50+
klog.ErrorS(errors.New("failed to assert ossfs mount result type"), "skipping monitoring of mountpoint", "mountpoint", req.Target)
5051
return err
5152
}
5253

@@ -65,6 +66,6 @@ func OssfsMonitorInterceptor(ctx context.Context, op *mounter.MountOperation, ha
6566

6667
monitor.HandleMountSuccess(res.PID)
6768
// Start monitoring goroutine (ticker based only)
68-
monitorManager.StartMonitoring(op.Target)
69+
monitorManager.StartMonitoring(req.Target)
6970
return nil
7071
}

0 commit comments

Comments
 (0)