• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

kubevirt / kubevirt / 9c9f5dca-1947-4c90-81dc-59a64b19e0eb

29 May 2026 09:13PM UTC coverage: 71.588% (+0.006%) from 71.582%
9c9f5dca-1947-4c90-81dc-59a64b19e0eb

push

prow

web-flow
Merge pull request #17959 from mhenriks/fix-vmexport-symlink-traversal

Fix symlink traversal in VMExport dir handler

17 of 19 new or added lines in 1 file covered. (89.47%)

78374 of 109480 relevant lines covered (71.59%)

565.99 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

55.5
/pkg/storage/export/virt-exportserver/exportserver.go
1
/*
2
 * This file is part of the KubeVirt project
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 *
16
 * Copyright The KubeVirt Authors.
17
 *
18
 */
19

20
package virtexportserver
21

22
import (
23
        "bytes"
24
        "context"
25
        "crypto/tls"
26
        "crypto/x509"
27
        "encoding/json"
28
        "errors"
29
        goflag "flag"
30
        "fmt"
31
        "io"
32
        golog "log"
33
        "net"
34
        "net/http"
35
        "os"
36
        "os/exec"
37
        "path"
38
        "path/filepath"
39
        "strconv"
40
        "strings"
41
        "sync"
42
        "time"
43

44
        gzip "github.com/klauspost/pgzip"
45
        flag "github.com/spf13/pflag"
46
        "golang.org/x/net/http2"
47
        "google.golang.org/grpc"
48
        "google.golang.org/grpc/credentials/insecure"
49
        "google.golang.org/grpc/keepalive"
50
        corev1 "k8s.io/api/core/v1"
51
        metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
52
        "k8s.io/apimachinery/pkg/runtime"
53
        "k8s.io/apimachinery/pkg/runtime/schema"
54
        "sigs.k8s.io/yaml"
55

56
        virtv1 "kubevirt.io/api/core/v1"
57
        "kubevirt.io/client-go/log"
58
        cdiv1 "kubevirt.io/containerized-data-importer-api/pkg/apis/core/v1beta1"
59

60
        nbdv1 "kubevirt.io/kubevirt/pkg/storage/cbt/nbd/v1"
61

62
        backupv1 "kubevirt.io/api/backup/v1alpha1"
63

64
        "kubevirt.io/kubevirt/pkg/safepath"
65
        "kubevirt.io/kubevirt/pkg/service"
66
        "kubevirt.io/kubevirt/pkg/storage/export/export"
67
        storageutils "kubevirt.io/kubevirt/pkg/storage/utils"
68
)
69

70
const (
71
        authHeader              = "x-kubevirt-export-token"
72
        manifestCmBasePath      = "/manifest_data/"
73
        vmManifestPath          = manifestCmBasePath + "virtualmachine-manifest"
74
        internalLinkPath        = manifestCmBasePath + "internal_host"
75
        internalCaConfigMapPath = manifestCmBasePath + "internal_ca_cm"
76
        externalLinkPath        = manifestCmBasePath + "external_host"
77
        externalCaConfigMapPath = manifestCmBasePath + "external_ca_cm"
78
        exportNamePath          = manifestCmBasePath + "export-name"
79

80
        external = "/external"
81
        internal = "/internal"
82

83
        defaultMapPageSize = 512
84
        tunnelIdleTimeout  = 60 * time.Second
85
)
86

87
var (
88
        excludeMap = map[string]struct{}{
89
                "lost+found": {},
90
        }
91
        h2DummyAddr = &net.TCPAddr{}
92
)
93

94
type TokenGetterFunc func() (string, error)
95

96
type ExportServerConfig struct {
97
        Deadline time.Time
98

99
        ListenAddr string
100

101
        CertFile, KeyFile string
102
        BackupCACert      []byte
103
        TLSMinVersion     uint16
104
        TLSCipherSuites   []uint16
105

106
        TokenFile string
107

108
        BackupUID        string
109
        BackupType       string
110
        BackupCheckpoint string
111

112
        Paths *export.ServerPaths
113

114
        // unit testing helpers
115
        ArchiveHandler     func(string) http.Handler
116
        DirHandler         func(string, string) http.Handler
117
        FileHandler        func(string) http.Handler
118
        GzipHandler        func(string) http.Handler
119
        VmHandler          func([]export.VolumeInfo, func() (string, error), func() (*corev1.ConfigMap, error)) http.Handler
120
        TokenSecretHandler func(TokenGetterFunc) http.Handler
121

122
        PermissionChecker func(string) bool
123

124
        TokenGetter TokenGetterFunc
125
}
126

127
type execReader struct {
128
        cmd    *exec.Cmd
129
        stdout io.ReadCloser
130
        stderr io.ReadCloser
131
}
132

133
type exportServer struct {
134
        ExportServerConfig
135
        handler http.Handler
136

137
        nbdClient nbdv1.NBDClient
138
        nbdMu     sync.RWMutex
139
}
140

141
func (er *execReader) Read(p []byte) (int, error) {
×
142
        n, err := er.stdout.Read(p)
×
143
        if err == io.EOF {
×
144
                if err2 := er.cmd.Wait(); err2 != nil {
×
145
                        errBytes, _ := io.ReadAll(er.stderr)
×
146
                        log.Log.Reason(err2).Errorf("Subprocess did not execute successfully, result is: %q\n%s", er.cmd.ProcessState.ExitCode(), string(errBytes))
×
147
                        return n, err2
×
148
                }
×
149
        }
150
        return n, err
×
151
}
152

153
func (er *execReader) Close() error {
×
154
        return er.stdout.Close()
×
155
}
×
156

157
func (s *exportServer) initHandler() {
24✔
158
        mux := http.NewServeMux()
24✔
159
        for _, vi := range s.Paths.Volumes {
40✔
160
                if hasPermissions := s.PermissionChecker(vi.Path); !hasPermissions {
16✔
161
                        golog.Fatalf("unable to manipulate %s's contents, exiting", vi.Path)
×
162
                }
×
163
                for path, handler := range s.getHandlerMap(vi) {
32✔
164
                        log.Log.Infof("Handling path %s\n", path)
16✔
165
                        mux.Handle(path, tokenChecker(s.TokenGetter, handler))
16✔
166
                }
16✔
167
        }
168
        for _, bi := range s.Paths.Backups {
24✔
169
                log.Log.Infof("Handling backup path %s (Map) and %s (Data)\n", bi.MapURI, bi.DataURI)
×
170
                mux.Handle(bi.MapURI, tokenChecker(s.TokenGetter, s.backupMapHandler(bi.Path)))
×
171
                mux.Handle(bi.DataURI, tokenChecker(s.TokenGetter, s.backupDataHandler(bi.Path)))
×
172
        }
×
173
        if s.Paths.VMURI != "" {
32✔
174
                mux.Handle(filepath.Join(internal, s.Paths.VMURI), tokenChecker(s.TokenGetter, s.VmHandler(s.Paths.Volumes, getInternalBasePath, getInternalCAConfigMap)))
8✔
175
                mux.Handle(filepath.Join(external, s.Paths.VMURI), tokenChecker(s.TokenGetter, s.VmHandler(s.Paths.Volumes, getExternalBasePath, getExternalCAConfigMap)))
8✔
176
        }
8✔
177
        if s.Paths.SecretURI != "" {
24✔
178
                mux.Handle(filepath.Join(internal, s.Paths.SecretURI), tokenChecker(s.TokenGetter, s.TokenSecretHandler(s.TokenGetter)))
×
179
                mux.Handle(filepath.Join(external, s.Paths.SecretURI), tokenChecker(s.TokenGetter, s.TokenSecretHandler(s.TokenGetter)))
×
180
        }
×
181
        // Readiness probe
182
        mux.HandleFunc(export.ReadinessPath, s.readyHandler)
24✔
183

24✔
184
        s.handler = mux
24✔
185
}
186

187
func getInternalCAConfigMap() (*corev1.ConfigMap, error) {
×
188
        return getCAConfigMap(internalCaConfigMapPath)
×
189
}
×
190

191
func getExternalCAConfigMap() (*corev1.ConfigMap, error) {
×
192
        return getCAConfigMap(externalCaConfigMapPath)
×
193
}
×
194

195
func (s *exportServer) getHandlerMap(vi export.VolumeInfo) map[string]http.Handler {
16✔
196
        fi, err := os.Stat(vi.Path)
16✔
197
        if err != nil {
16✔
198
                log.Log.Reason(err).Errorf("error statting %s", vi.Path)
×
199
                return nil
×
200
        }
×
201

202
        var result = make(map[string]http.Handler)
16✔
203

16✔
204
        if vi.ArchiveURI != "" {
20✔
205
                result[vi.ArchiveURI] = s.ArchiveHandler(vi.Path)
4✔
206
        }
4✔
207

208
        if vi.DirURI != "" {
20✔
209
                result[vi.DirURI] = s.DirHandler(vi.DirURI, vi.Path)
4✔
210
        }
4✔
211

212
        p := vi.Path
16✔
213
        if fi.IsDir() {
32✔
214
                p = path.Join(p, "disk.img")
16✔
215
        }
16✔
216

217
        if vi.RawURI != "" {
20✔
218
                result[vi.RawURI] = s.FileHandler(p)
4✔
219
        }
4✔
220

221
        if vi.RawGzURI != "" {
20✔
222
                result[vi.RawGzURI] = s.GzipHandler(p)
4✔
223
        }
4✔
224

225
        return result
16✔
226
}
227

228
func (s *exportServer) Run() {
×
229
        s.initHandler()
×
230

×
231
        srv := s.buildServer()
×
232

×
233
        h2Server := &http2.Server{
×
234
                IdleTimeout: tunnelIdleTimeout,
×
235
        }
×
236
        if err := http2.ConfigureServer(srv, h2Server); err != nil {
×
237
                panic(err)
×
238
        }
239

240
        ch := make(chan error)
×
241

×
242
        go func() {
×
243
                err := srv.ListenAndServeTLS(s.CertFile, s.KeyFile)
×
244
                ch <- err
×
245
        }()
×
246

247
        if !s.Deadline.IsZero() {
×
248
                log.Log.Infof("Deadline set to %s", s.Deadline)
×
249
                select {
×
250
                case err := <-ch:
×
251
                        panic(err)
×
252
                case <-time.After(time.Until(s.Deadline)):
×
253
                        log.Log.Info("Deadline exceeded, shutting down")
×
254
                        srv.Shutdown(context.TODO())
×
255
                }
256
        } else {
×
257
                err := <-ch
×
258
                panic(err)
×
259
        }
260
}
261

262
func (s *exportServer) buildServer() *http.Server {
3✔
263
        tlsConfig := &tls.Config{
3✔
264
                MinVersion:   s.TLSMinVersion,
3✔
265
                CipherSuites: s.TLSCipherSuites,
3✔
266
                NextProtos:   []string{"h2", "http/1.1"},
3✔
267
        }
3✔
268

3✔
269
        rootHandler := s.handler
3✔
270
        if s.BackupUID != "" {
5✔
271
                clientCAPool := x509.NewCertPool()
2✔
272
                if ok := clientCAPool.AppendCertsFromPEM(s.BackupCACert); !ok {
3✔
273
                        panic("failed to parse Backup CA")
1✔
274
                }
275
                tlsConfig.ClientCAs = clientCAPool
1✔
276
                tlsConfig.ClientAuth = tls.VerifyClientCertIfGiven
1✔
277

1✔
278
                rootHandler = http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
3✔
279
                        if r.Method == http.MethodConnect {
3✔
280
                                s.handleTunnel(w, r)
1✔
281
                                return
1✔
282
                        }
1✔
283
                        s.handler.ServeHTTP(w, r)
1✔
284
                })
285
        }
286

287
        return &http.Server{
2✔
288
                Addr:      s.ListenAddr,
2✔
289
                Handler:   rootHandler,
2✔
290
                TLSConfig: tlsConfig,
2✔
291
        }
2✔
292
}
293

294
func (s *exportServer) AddFlags() {
×
295
        flag.CommandLine.AddGoFlag(goflag.CommandLine.Lookup("v"))
×
296
}
×
297

298
func NewExportServer(config ExportServerConfig) service.Service {
24✔
299
        es := &exportServer{
24✔
300
                ExportServerConfig: config,
24✔
301
        }
24✔
302

24✔
303
        if es.ArchiveHandler == nil {
24✔
304
                es.ArchiveHandler = archiveHandler
×
305
        }
×
306

307
        if es.DirHandler == nil {
24✔
308
                es.DirHandler = dirHandler
×
309
        }
×
310

311
        if es.FileHandler == nil {
24✔
312
                es.FileHandler = fileHandler
×
313
        }
×
314

315
        if es.GzipHandler == nil {
24✔
316
                es.GzipHandler = gzipHandler
×
317
        }
×
318

319
        if es.VmHandler == nil {
24✔
320
                es.VmHandler = vmHandler
×
321
        }
×
322

323
        if es.TokenSecretHandler == nil {
24✔
324
                es.TokenSecretHandler = secretHandler
×
325
        }
×
326

327
        if es.TokenGetter == nil {
24✔
328
                es.TokenGetter = func() (string, error) {
×
329
                        return getToken(es.TokenFile)
×
330
                }
×
331
        }
332

333
        if es.PermissionChecker == nil {
24✔
334
                es.PermissionChecker = checkVolumePermissions
×
335
        }
×
336

337
        return es
24✔
338
}
339

340
var getExpandedVM = func() *virtv1.VirtualMachine {
×
341
        f, err := os.Open(vmManifestPath)
×
342
        if err != nil {
×
343
                log.Log.Reason(err).Info("Unable to load VM manifest data")
×
344
                return nil
×
345
        }
×
346
        defer f.Close()
×
347
        fileinfo, err := f.Stat()
×
348
        if err != nil {
×
349
                log.Log.Reason(err).Info("Unable to load VM manifest data")
×
350
                return nil
×
351
        }
×
352
        buf := make([]byte, fileinfo.Size())
×
353
        _, err = f.Read(buf)
×
354
        if err != nil {
×
355
                log.Log.Reason(err).Info("Unable to load VM manifest data")
×
356
                return nil
×
357
        }
×
358

359
        vm := &virtv1.VirtualMachine{}
×
360
        if err := json.Unmarshal(buf, vm); err != nil {
×
361
                log.Log.Reason(err).Info("Unable to load VM manifest data")
×
362
                return nil
×
363
        }
×
364
        return vm
×
365
}
366

367
var getInternalBasePath = func() (string, error) {
×
368
        data, err := os.ReadFile(internalLinkPath)
×
369
        if err != nil {
×
370
                return "", err
×
371
        }
×
372
        return string(data), nil
×
373
}
374

375
var getExportName = func() (string, error) {
×
376
        data, err := os.ReadFile(exportNamePath)
×
377
        if err != nil {
×
378
                return "", err
×
379
        }
×
380
        return string(data), nil
×
381
}
382

383
var getExternalBasePath = func() (string, error) {
×
384
        data, err := os.ReadFile(externalLinkPath)
×
385
        if err != nil {
×
386
                return "", err
×
387
        }
×
388
        return string(data), nil
×
389
}
390

391
func GetTypeMetaString(gvk schema.GroupVersionKind) string {
×
392
        return fmt.Sprintf("apiVersion: %s\nkind: %s\n", gvk.GroupVersion().String(), gvk.Kind)
×
393
}
×
394

395
var getCAConfigMap = func(name string) (*corev1.ConfigMap, error) {
×
396
        f, err := os.Open(name)
×
397
        if err != nil {
×
398
                return nil, err
×
399
        }
×
400
        defer f.Close()
×
401
        fileinfo, err := f.Stat()
×
402
        if err != nil {
×
403
                return nil, err
×
404
        }
×
405
        buf := make([]byte, fileinfo.Size())
×
406
        _, err = f.Read(buf)
×
407
        if err != nil {
×
408
                return nil, err
×
409
        }
×
410

411
        cm := &corev1.ConfigMap{}
×
412
        if err := json.Unmarshal(buf, cm); err != nil {
×
413
                return nil, err
×
414
        }
×
415
        return cm, nil
×
416
}
417

418
var getCdiHeaderSecret = func(token, name string) *corev1.Secret {
2✔
419
        data := make(map[string]string)
2✔
420

2✔
421
        data["token"] = fmt.Sprintf("x-kubevirt-export-token:%s", token)
2✔
422
        return &corev1.Secret{
2✔
423
                ObjectMeta: metav1.ObjectMeta{
2✔
424
                        Name: name,
2✔
425
                },
2✔
426
                StringData: data,
2✔
427
        }
2✔
428
}
2✔
429

430
var getDataVolumes = func(vm *virtv1.VirtualMachine) ([]*cdiv1.DataVolume, error) {
×
431
        res := make([]*cdiv1.DataVolume, 0)
×
432
        volumes, err := storageutils.GetVolumes(vm, nil)
×
433
        if err != nil {
×
434
                return nil, err
×
435
        }
×
436
        for _, volume := range volumes {
×
437
                name := ""
×
438
                if volume.DataVolume != nil {
×
439
                        name = volume.DataVolume.Name
×
440
                } else if volume.PersistentVolumeClaim != nil {
×
441
                        name = volume.PersistentVolumeClaim.ClaimName
×
442
                }
×
443
                if name == "" {
×
444
                        continue
×
445
                }
446
                log.Log.V(1).Infof("Opening DV %s", filepath.Join(manifestCmBasePath, fmt.Sprintf("dv-%s", name)))
×
447
                f, err := os.Open(filepath.Join(manifestCmBasePath, fmt.Sprintf("dv-%s", name)))
×
448
                if err != nil {
×
449
                        if errors.Is(err, os.ErrNotExist) {
×
450
                                log.Log.V(1).Info("DV not found skipping")
×
451
                                continue
×
452
                        }
453
                        return nil, err
×
454
                }
455
                defer f.Close()
×
456
                fileinfo, err := f.Stat()
×
457
                if err != nil {
×
458
                        return nil, err
×
459
                }
×
460
                buf := make([]byte, fileinfo.Size())
×
461
                _, err = f.Read(buf)
×
462
                if err != nil {
×
463
                        return nil, err
×
464
                }
×
465
                dv := &cdiv1.DataVolume{}
×
466
                if err := json.Unmarshal(buf, dv); err != nil {
×
467
                        return nil, err
×
468
                }
×
469
                res = append(res, dv)
×
470
        }
471
        return res, nil
×
472
}
473

474
func newTarReader(mountPoint string) (io.ReadCloser, error) {
×
475
        var excludeArgs []string
×
476
        for name := range excludeMap {
×
477
                excludeArgs = append(excludeArgs, "--exclude="+name)
×
478
        }
×
479

480
        args := []string{"Scv"}
×
481
        args = append(args, excludeArgs...)
×
482
        args = append(args, ".")
×
483

×
484
        cmd := exec.Command("/usr/bin/tar", args...)
×
485
        cmd.Dir = mountPoint
×
486
        stdout, err := cmd.StdoutPipe()
×
487
        if err != nil {
×
488
                return nil, err
×
489
        }
×
490
        var stderr bytes.Buffer
×
491
        cmd.Stderr = &stderr
×
492
        if err = cmd.Start(); err != nil {
×
493
                return nil, err
×
494
        }
×
495
        return &execReader{cmd: cmd, stdout: stdout, stderr: io.NopCloser(&stderr)}, nil
×
496
}
497

498
func pipeToGzip(reader io.ReadCloser) io.ReadCloser {
×
499
        pr, pw := io.Pipe()
×
500
        zw := gzip.NewWriter(pw)
×
501

×
502
        go func() {
×
503
                n, err := io.Copy(zw, reader)
×
504
                if err != nil {
×
505
                        log.Log.Reason(err).Error("error piping to gzip")
×
506
                }
×
507
                if err = zw.Close(); err != nil {
×
508
                        log.Log.Reason(err).Error("error closing gzip writer")
×
509
                }
×
510
                if err = pw.Close(); err != nil {
×
511
                        log.Log.Reason(err).Error("error closing pipe writer")
×
512
                }
×
513
                log.Log.Infof("Wrote %d bytes\n", n)
×
514
        }()
515

516
        return pr
×
517
}
518

519
func getTokenQueryParam(r *http.Request) (token string) {
24✔
520
        q := r.URL.Query()
24✔
521
        if keys, ok := q[authHeader]; ok {
36✔
522
                token = keys[0]
12✔
523
                q.Del(authHeader)
12✔
524
                r.URL.RawQuery = q.Encode()
12✔
525
        }
12✔
526
        return
24✔
527
}
528

529
func getTokenHeader(r *http.Request) (token string) {
24✔
530
        if tok := r.Header.Get(authHeader); tok != "" {
36✔
531
                r.Header.Del(authHeader)
12✔
532
                token = tok
12✔
533
        }
12✔
534
        return
24✔
535
}
536

537
func tokenChecker(tokenGetter TokenGetterFunc, nextHandler http.Handler) http.Handler {
32✔
538
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
56✔
539
                token, err := tokenGetter()
24✔
540
                if err != nil {
24✔
541
                        w.WriteHeader(http.StatusInternalServerError)
×
542
                        return
×
543
                }
×
544
                for _, tok := range []string{getTokenQueryParam(r), getTokenHeader(r)} {
66✔
545
                        if tok == token {
54✔
546
                                nextHandler.ServeHTTP(w, r)
12✔
547
                                return
12✔
548
                        }
12✔
549
                }
550
                w.WriteHeader(http.StatusUnauthorized)
12✔
551
        })
552
}
553

554
func archiveHandler(mountPoint string) http.Handler {
×
555
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
×
556
                if req.Method != http.MethodGet {
×
557
                        w.WriteHeader(http.StatusBadRequest)
×
558
                        return
×
559
                }
×
560
                if hasPermissions := checkDirectoryPermissions(mountPoint); !hasPermissions {
×
561
                        w.WriteHeader(http.StatusInternalServerError)
×
562
                        return
×
563
                }
×
564

565
                tarReader, err := newTarReader(mountPoint)
×
566
                if err != nil {
×
567
                        log.Log.Reason(err).Error("error creating tar reader")
×
568
                        w.WriteHeader(http.StatusInternalServerError)
×
569
                        return
×
570
                }
×
571
                defer tarReader.Close()
×
572
                gzipReader := pipeToGzip(tarReader)
×
573
                defer gzipReader.Close()
×
574
                n, err := io.Copy(w, gzipReader)
×
575
                if err != nil {
×
576
                        log.Log.Reason(err).Error("error writing response body")
×
577
                }
×
578
                log.Log.Infof("Wrote %d bytes\n", n)
×
579
        })
580
}
581

582
func checkDirectoryPermissions(filePath string) bool {
×
583
        dir, err := os.Open(filePath)
×
584
        if err != nil {
×
585
                log.Log.Reason(err).Errorf("error opening %s", filePath)
×
586
                return false
×
587
        }
×
588
        defer dir.Close()
×
589

×
590
        // Read all filenames
×
591
        contents, err := dir.Readdirnames(-1)
×
592
        if err != nil {
×
593
                log.Log.Reason(err).Errorf("failed to read directory contents: %v", err)
×
594
                return false
×
595
        }
×
596

597
        for _, item := range contents {
×
598
                if _, ok := excludeMap[item]; ok {
×
599
                        continue
×
600
                }
601
                itemPath := filepath.Join(filePath, item)
×
602
                // Check if export server has permissions to manipulate the file
×
603
                file, err := os.Open(itemPath)
×
604
                if err != nil {
×
605
                        log.Log.Reason(err).Errorf("%s may lack read permissions", itemPath)
×
606
                        return false
×
607
                }
×
608
                file.Close()
×
609
        }
610
        return true
×
611
}
612

613
func checkVolumePermissions(path string) bool {
×
614
        fi, err := os.Stat(path)
×
615
        if err != nil {
×
616
                log.Log.Reason(err).Errorf("error statting %s", path)
×
617
                return false
×
618
        }
×
619
        if fi.IsDir() {
×
620
                return checkDirectoryPermissions(path)
×
621
        }
×
622
        f, err := os.Open(path)
×
623
        if err != nil {
×
624
                log.Log.Reason(err).Errorf("error opening %s", path)
×
625
                return false
×
626
        }
×
627
        f.Close()
×
628
        return true
×
629
}
630

631
func gzipHandler(filePath string) http.Handler {
×
632
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
×
633
                if req.Method != http.MethodGet {
×
634
                        w.WriteHeader(http.StatusBadRequest)
×
635
                        return
×
636
                }
×
637
                f, err := os.Open(filePath)
×
638
                if err != nil {
×
639
                        log.Log.Reason(err).Errorf("error opening %s", filePath)
×
640
                        w.WriteHeader(http.StatusInternalServerError)
×
641
                        return
×
642
                }
×
643
                defer f.Close()
×
644
                gzipReader := pipeToGzip(f)
×
645
                defer gzipReader.Close()
×
646
                n, err := io.Copy(w, gzipReader)
×
647
                if err != nil {
×
648
                        log.Log.Reason(err).Error("error writing response body")
×
649
                }
×
650
                log.Log.Infof("Wrote %d bytes\n", n)
×
651
        })
652
}
653

654
func vmHandler(vi []export.VolumeInfo, getBasePath func() (string, error), getCmFunc func() (*corev1.ConfigMap, error)) http.Handler {
14✔
655
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
28✔
656
                if req.Method != http.MethodGet {
18✔
657
                        w.WriteHeader(http.StatusBadRequest)
4✔
658
                        return
4✔
659
                }
4✔
660
                resources := make([]runtime.Object, 0)
10✔
661
                outputFunc := resourceToBytesJson
10✔
662
                contentType := req.Header.Get("Accept")
10✔
663
                if contentType == runtime.ContentTypeYAML {
13✔
664
                        outputFunc = resourceToBytesYaml
3✔
665
                }
3✔
666
                exportName, err := getExportName()
10✔
667
                if err != nil {
11✔
668
                        log.Log.Reason(err).Error("error reading export name")
1✔
669
                        w.WriteHeader(http.StatusInternalServerError)
1✔
670
                        return
1✔
671
                }
1✔
672
                headerSecretName := getSecretTokenName(exportName)
9✔
673
                path, err := getBasePath()
9✔
674
                if err != nil {
11✔
675
                        if errors.Is(err, os.ErrNotExist) {
2✔
676
                                log.Log.Reason(err).Info("path not found")
×
677
                                w.WriteHeader(http.StatusNotFound)
×
678
                        } else {
2✔
679
                                log.Log.Reason(err).Error("error reading path")
2✔
680
                                w.WriteHeader(http.StatusInternalServerError)
2✔
681
                        }
2✔
682
                        return
2✔
683
                }
684
                certCm, error := getCmFunc()
7✔
685
                if error != nil {
8✔
686
                        log.Log.Reason(err).Error("error reading ca configmap information")
1✔
687
                        w.WriteHeader(http.StatusInternalServerError)
1✔
688
                        return
1✔
689
                }
1✔
690
                certCm.TypeMeta = metav1.TypeMeta{
6✔
691
                        Kind:       "ConfigMap",
6✔
692
                        APIVersion: "v1",
6✔
693
                }
6✔
694
                resources = append(resources, certCm)
6✔
695
                expandedVm := getExpandedVM()
6✔
696
                if expandedVm == nil {
7✔
697
                        log.Log.Reason(err).Error("error getting VM definition")
1✔
698
                        w.WriteHeader(http.StatusInternalServerError)
1✔
699
                        return
1✔
700
                }
1✔
701
                expandedVm.TypeMeta = metav1.TypeMeta{
5✔
702
                        Kind:       virtv1.VirtualMachineGroupVersionKind.Kind,
5✔
703
                        APIVersion: virtv1.VirtualMachineGroupVersionKind.GroupVersion().String(),
5✔
704
                }
5✔
705
                for i, dvTemplate := range expandedVm.Spec.DataVolumeTemplates {
7✔
706
                        dvTemplate.Spec.Source.HTTP.URL = fmt.Sprintf("https://%s", filepath.Join(path, vi[i].RawGzURI))
2✔
707
                        dvTemplate.Spec.Source.HTTP.CertConfigMap = certCm.Name
2✔
708
                        dvTemplate.Spec.Source.HTTP.SecretExtraHeaders = []string{headerSecretName}
2✔
709
                }
2✔
710
                resources = append(resources, expandedVm)
5✔
711
                datavolumes, err := getDataVolumes(expandedVm)
5✔
712
                if err != nil {
5✔
713
                        log.Log.Reason(err).Error("error reading datavolumes information")
×
714
                        w.WriteHeader(http.StatusInternalServerError)
×
715
                        return
×
716
                }
×
717
                for _, dv := range datavolumes {
6✔
718
                        dv.TypeMeta = metav1.TypeMeta{
1✔
719
                                Kind:       "DataVolume",
1✔
720
                                APIVersion: "cdi.kubevirt.io/v1beta1",
1✔
721
                        }
1✔
722
                        for _, info := range vi {
2✔
723
                                if strings.Contains(info.RawGzURI, dv.Name) {
2✔
724
                                        dv.Spec.Source.HTTP.URL = fmt.Sprintf("https://%s", filepath.Join(path, info.RawGzURI))
1✔
725
                                }
1✔
726
                        }
727
                        dv.Spec.Source.HTTP.CertConfigMap = certCm.Name
1✔
728
                        dv.Spec.Source.HTTP.SecretExtraHeaders = []string{headerSecretName}
1✔
729
                        resources = append(resources, dv)
1✔
730
                }
731
                data, err := outputFunc(resources)
5✔
732
                if err != nil {
5✔
733
                        w.WriteHeader(http.StatusInternalServerError)
×
734
                        return
×
735
                }
×
736
                n, err := w.Write(data)
5✔
737
                if err != nil {
5✔
738
                        log.Log.Reason(err).Error("error writing manifests")
×
739
                        w.WriteHeader(http.StatusInternalServerError)
×
740
                        return
×
741
                }
×
742
                log.Log.Infof("Wrote %d bytes\n", n)
5✔
743
        })
744
}
745

746
func resourceToBytesJson(resources []runtime.Object) ([]byte, error) {
3✔
747
        list := corev1.List{
3✔
748
                TypeMeta: metav1.TypeMeta{
3✔
749
                        Kind:       "List",
3✔
750
                        APIVersion: "v1",
3✔
751
                },
3✔
752
                ListMeta: metav1.ListMeta{},
3✔
753
        }
3✔
754
        for _, resource := range resources {
8✔
755
                list.Items = append(list.Items, runtime.RawExtension{Object: resource})
5✔
756
        }
5✔
757
        resourceBytes, err := json.MarshalIndent(list, "", "    ")
3✔
758
        if err != nil {
3✔
759
                return nil, err
×
760
        }
×
761
        return resourceBytes, nil
3✔
762
}
763

764
func resourceToBytesYaml(resources []runtime.Object) ([]byte, error) {
4✔
765
        data := []byte{}
4✔
766
        for _, resource := range resources {
12✔
767
                resourceBytes, err := yaml.Marshal(resource)
8✔
768
                if err != nil {
8✔
769
                        return nil, err
×
770
                }
×
771
                data = append(data, resourceBytes...)
8✔
772
                data = append(data, []byte("---\n")...)
8✔
773
        }
774
        return data, nil
4✔
775
}
776

777
// symlinkSafeDir is an http.FileSystem that prevents symlink traversal outside root.
778
// Unlike http.Dir, it resolves symlinks via safepath and rejects any that escape the
779
// root boundary, preventing attackers from reading files outside the exported PVC.
780
type symlinkSafeDir struct {
781
        root string
782
}
783

784
func (d symlinkSafeDir) Open(name string) (http.File, error) {
7✔
785
        cleanName := path.Clean("/" + name)
7✔
786
        if cleanName == "/" {
8✔
787
                cleanName = "."
1✔
788
        } else {
7✔
789
                cleanName = cleanName[1:]
6✔
790
        }
6✔
791

792
        resolved, err := safepath.JoinAndResolveWithRelativeRoot(d.root, cleanName)
7✔
793
        if err != nil {
11✔
794
                return nil, err
4✔
795
        }
4✔
796

797
        fd, err := safepath.OpenAtNoFollow(resolved)
3✔
798
        if err != nil {
3✔
NEW
799
                return nil, err
×
NEW
800
        }
×
801
        defer fd.Close()
3✔
802

3✔
803
        return os.Open(fd.SafePath())
3✔
804
}
805

806
func dirHandler(uri, mountPoint string) http.Handler {
6✔
807
        return http.StripPrefix(uri, http.FileServer(symlinkSafeDir{root: mountPoint}))
6✔
808
}
6✔
809

810
func fileHandler(file string) http.Handler {
×
811
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
×
812
                f, err := os.Open(file)
×
813
                if err != nil {
×
814
                        log.Log.Reason(err).Errorf("error opening %s", file)
×
815
                        w.WriteHeader(http.StatusInternalServerError)
×
816
                        return
×
817
                }
×
818
                defer f.Close()
×
819
                http.ServeContent(w, r, "disk.img", time.Time{}, f)
×
820
        })
821
}
822

823
func getToken(tokenFile string) (string, error) {
×
824
        content, err := os.ReadFile(tokenFile)
×
825
        if err != nil {
×
826
                return "", err
×
827
        }
×
828

829
        return string(content), nil
×
830
}
831

832
var getSecretTokenName = func(exportName string) string {
11✔
833
        return fmt.Sprintf("header-secret-%s", exportName)
11✔
834
}
11✔
835

836
func secretHandler(tokenGetter TokenGetterFunc) http.Handler {
8✔
837
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
16✔
838
                if req.Method != http.MethodGet {
12✔
839
                        w.WriteHeader(http.StatusBadRequest)
4✔
840
                        return
4✔
841
                }
4✔
842
                resources := make([]runtime.Object, 0)
4✔
843
                outputFunc := resourceToBytesJson
4✔
844
                contentType := req.Header.Get("Accept")
4✔
845
                if contentType == runtime.ContentTypeYAML {
5✔
846
                        outputFunc = resourceToBytesYaml
1✔
847
                }
1✔
848
                token, err := tokenGetter()
4✔
849
                if err != nil {
5✔
850
                        log.Log.Reason(err).Error("error getting token")
1✔
851
                        w.WriteHeader(http.StatusInternalServerError)
1✔
852
                        return
1✔
853
                }
1✔
854
                exportName, err := getExportName()
3✔
855
                if err != nil {
4✔
856
                        log.Log.Reason(err).Error("error reading export name")
1✔
857
                        w.WriteHeader(http.StatusInternalServerError)
1✔
858
                        return
1✔
859
                }
1✔
860
                headerSecretName := getSecretTokenName(exportName)
2✔
861
                secret := getCdiHeaderSecret(token, headerSecretName)
2✔
862
                secret.TypeMeta = metav1.TypeMeta{
2✔
863
                        Kind:       "Secret",
2✔
864
                        APIVersion: "v1",
2✔
865
                }
2✔
866
                resources = append(resources, secret)
2✔
867
                data, err := outputFunc(resources)
2✔
868
                if err != nil {
2✔
869
                        log.Log.Reason(err).Errorf("error generating secret manifest")
×
870
                        w.WriteHeader(http.StatusInternalServerError)
×
871
                        return
×
872
                }
×
873
                n, err := w.Write(data)
2✔
874
                if err != nil {
2✔
875
                        log.Log.Reason(err).Error("error writing secret manifest")
×
876
                        w.WriteHeader(http.StatusInternalServerError)
×
877
                        return
×
878
                }
×
879
                log.Log.Infof("Wrote %d bytes\n", n)
2✔
880
        })
881
}
882

883
func (s *exportServer) readyHandler(w http.ResponseWriter, r *http.Request) {
×
884
        io.WriteString(w, "OK")
×
885
}
×
886

887
func (s *exportServer) handleTunnel(w http.ResponseWriter, r *http.Request) {
5✔
888
        if r.TLS == nil || len(r.TLS.PeerCertificates) == 0 {
7✔
889
                log.Log.Error("tunnel rejected: no client certificate presented")
2✔
890
                http.Error(w, "mTLS required", http.StatusUnauthorized)
2✔
891
                return
2✔
892
        }
2✔
893

894
        expectedCN := fmt.Sprintf("kubevirt.io:system:client:%s", s.BackupUID)
3✔
895
        clientCN := r.TLS.PeerCertificates[0].Subject.CommonName
3✔
896
        if clientCN != expectedCN {
4✔
897
                log.Log.Errorf("identity mismatch, cert: %s, expected: %s", clientCN, expectedCN)
1✔
898
                http.Error(w, "Forbidden", http.StatusForbidden)
1✔
899
                return
1✔
900
        }
1✔
901

902
        s.nbdMu.Lock()
2✔
903
        if s.nbdClient != nil {
3✔
904
                s.nbdMu.Unlock()
1✔
905
                _ = r.Body.Close()
1✔
906
                log.Log.Warning("rejecting tunnel: active session already exists")
1✔
907
                http.Error(w, "Conflict", http.StatusConflict)
1✔
908
                return
1✔
909
        }
1✔
910

911
        ctx, cancel := context.WithCancel(r.Context())
1✔
912
        defer cancel()
1✔
913
        conn := newH2ServerConn(r.Body, w, cancel)
1✔
914

1✔
915
        var dialOnce sync.Once
1✔
916
        clientConn, err := grpc.NewClient(
1✔
917
                "passthrough:///backup",
1✔
918
                grpc.WithContextDialer(func(_ context.Context, _ string) (net.Conn, error) {
1✔
919
                        var c net.Conn
×
920
                        dialOnce.Do(func() { c = conn })
×
921
                        if c != nil {
×
922
                                return c, nil
×
923
                        }
×
924
                        return nil, fmt.Errorf("tunnel connection is single-use; reconnect not supported")
×
925
                }),
926
                grpc.WithTransportCredentials(insecure.NewCredentials()),
927
                grpc.WithKeepaliveParams(keepalive.ClientParameters{
928
                        Time:                10 * time.Second,
929
                        Timeout:             5 * time.Second,
930
                        PermitWithoutStream: true,
931
                }),
932
        )
933
        if err != nil {
1✔
934
                s.nbdMu.Unlock()
×
935
                log.Log.Reason(err).Error("failed to initialize gRPC client for tunnel")
×
936
                http.Error(w, "Internal Server Error", http.StatusInternalServerError)
×
937
                return
×
938
        }
×
939

940
        w.WriteHeader(http.StatusOK)
1✔
941
        w.(http.Flusher).Flush()
1✔
942

1✔
943
        s.nbdClient = nbdv1.NewNBDClient(clientConn)
1✔
944
        s.nbdMu.Unlock()
1✔
945

1✔
946
        log.Log.Infof("Exclusive backup tunnel established for %s", s.BackupUID)
1✔
947

1✔
948
        <-ctx.Done()
1✔
949

1✔
950
        s.nbdMu.Lock()
1✔
951
        clientConn.Close()
1✔
952
        s.nbdClient = nil
1✔
953
        s.nbdMu.Unlock()
1✔
954
        log.Log.Info("Backup tunnel disconnected, listener reset")
1✔
955
}
956

957
type h2ServerConn struct {
958
        r      io.ReadCloser
959
        w      http.ResponseWriter
960
        cancel context.CancelFunc
961
        once   sync.Once
962
}
963

964
func newH2ServerConn(r io.ReadCloser, w http.ResponseWriter, cancel context.CancelFunc) *h2ServerConn {
2✔
965
        return &h2ServerConn{r: r, w: w, cancel: cancel}
2✔
966
}
2✔
967

968
func (c *h2ServerConn) Read(b []byte) (int, error) { return c.r.Read(b) }
1✔
969

970
func (c *h2ServerConn) Write(b []byte) (int, error) {
1✔
971
        n, err := c.w.Write(b)
1✔
972
        if f, ok := c.w.(http.Flusher); ok {
2✔
973
                f.Flush()
1✔
974
        }
1✔
975
        return n, err
1✔
976
}
977

978
func (c *h2ServerConn) Close() error {
1✔
979
        c.once.Do(c.cancel)
1✔
980
        return c.r.Close()
1✔
981
}
1✔
982

983
func (c *h2ServerConn) LocalAddr() net.Addr                { return h2DummyAddr }
×
984
func (c *h2ServerConn) RemoteAddr() net.Addr               { return h2DummyAddr }
×
985
func (c *h2ServerConn) SetDeadline(_ time.Time) error      { return nil }
×
986
func (c *h2ServerConn) SetReadDeadline(_ time.Time) error  { return nil }
×
987
func (c *h2ServerConn) SetWriteDeadline(_ time.Time) error { return nil }
×
988

989
type ExportMapExtent struct {
990
        Offset      uint64 `json:"offset"`
991
        Length      uint64 `json:"length"`
992
        Type        uint64 `json:"type"`
993
        Description string `json:"description"`
994
}
995

996
type ExportMapResponse struct {
997
        Extents    []ExportMapExtent `json:"extents"`
998
        NextOffset *uint64           `json:"next_offset"`
999
}
1000

1001
func (s *exportServer) backupMapHandler(exportName string) http.Handler {
15✔
1002
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
30✔
1003
                if req.Method != http.MethodGet {
19✔
1004
                        w.WriteHeader(http.StatusMethodNotAllowed)
4✔
1005
                        return
4✔
1006
                }
4✔
1007

1008
                s.nbdMu.RLock()
11✔
1009
                client := s.nbdClient
11✔
1010
                s.nbdMu.RUnlock()
11✔
1011
                if client == nil {
12✔
1012
                        http.Error(w, "Backup source (virt-launcher) not connected via tunnel", http.StatusServiceUnavailable)
1✔
1013
                        return
1✔
1014
                }
1✔
1015

1016
                offset := uint64(0)
10✔
1017
                length := uint64(0)
10✔
1018
                pageSize := defaultMapPageSize
10✔
1019
                query := req.URL.Query()
10✔
1020

10✔
1021
                if offsetStr := query.Get("offset"); offsetStr != "" {
12✔
1022
                        o, err := strconv.ParseUint(offsetStr, 10, 64)
2✔
1023
                        if err != nil {
3✔
1024
                                http.Error(w, fmt.Sprintf("invalid offset %q: %v", offsetStr, err), http.StatusBadRequest)
1✔
1025
                                return
1✔
1026
                        }
1✔
1027
                        offset = o
1✔
1028
                }
1029
                if lengthStr := query.Get("length"); lengthStr != "" {
11✔
1030
                        l, err := strconv.ParseUint(lengthStr, 10, 64)
2✔
1031
                        if err != nil {
3✔
1032
                                http.Error(w, fmt.Sprintf("invalid length %q: %v", lengthStr, err), http.StatusBadRequest)
1✔
1033
                                return
1✔
1034
                        }
1✔
1035
                        length = l
1✔
1036
                }
1037
                if pageSizeStr := query.Get("page_size"); pageSizeStr != "" {
11✔
1038
                        p, err := strconv.Atoi(pageSizeStr)
3✔
1039
                        if err != nil || p <= 0 {
5✔
1040
                                http.Error(w, fmt.Sprintf("invalid page_size %q", pageSizeStr), http.StatusBadRequest)
2✔
1041
                                return
2✔
1042
                        }
2✔
1043
                        pageSize = p
1✔
1044
                }
1045

1046
                var bitmapName string
6✔
1047
                if s.BackupType == string(backupv1.Incremental) && s.BackupCheckpoint != "" {
7✔
1048
                        bitmapName = s.BackupCheckpoint
1✔
1049
                }
1✔
1050

1051
                streamCtx, streamCancel := context.WithCancel(req.Context())
6✔
1052
                defer streamCancel()
6✔
1053

6✔
1054
                stream, err := client.Map(streamCtx, &nbdv1.MapRequest{
6✔
1055
                        ExportName: exportName,
6✔
1056
                        BitmapName: bitmapName,
6✔
1057
                        Offset:     offset,
6✔
1058
                        Length:     length,
6✔
1059
                })
6✔
1060
                if err != nil {
7✔
1061
                        errMsg := fmt.Sprintf("Failed to call map for export: %s", exportName)
1✔
1062
                        log.Log.Reason(err).Error(errMsg)
1✔
1063
                        http.Error(w, errMsg, http.StatusInternalServerError)
1✔
1064
                        return
1✔
1065
                }
1✔
1066

1067
                extents, nextOffsetPtr, err := collectMapPage(stream, pageSize)
5✔
1068
                if err != nil {
6✔
1069
                        errMsg := fmt.Sprintf("Failed to collect map extents for export: %s", exportName)
1✔
1070
                        log.Log.Reason(err).Error(errMsg)
1✔
1071
                        http.Error(w, errMsg, http.StatusInternalServerError)
1✔
1072
                        return
1✔
1073
                }
1✔
1074

1075
                page := ExportMapResponse{
4✔
1076
                        Extents:    extents,
4✔
1077
                        NextOffset: nextOffsetPtr,
4✔
1078
                }
4✔
1079

4✔
1080
                w.Header().Set("Content-Type", "application/json")
4✔
1081
                if err := json.NewEncoder(w).Encode(page); err != nil {
4✔
1082
                        log.Log.Reason(err).Errorf("failed to encode map page for export %s", exportName)
×
1083
                }
×
1084
        })
1085
}
1086

1087
func (s *exportServer) backupDataHandler(exportName string) http.Handler {
11✔
1088
        return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
22✔
1089
                if req.Method != http.MethodGet {
15✔
1090
                        w.WriteHeader(http.StatusMethodNotAllowed)
4✔
1091
                        return
4✔
1092
                }
4✔
1093

1094
                s.nbdMu.RLock()
7✔
1095
                client := s.nbdClient
7✔
1096
                s.nbdMu.RUnlock()
7✔
1097

7✔
1098
                if client == nil {
8✔
1099
                        http.Error(w, "Backup source not connected", http.StatusServiceUnavailable)
1✔
1100
                        return
1✔
1101
                }
1✔
1102

1103
                offset := uint64(0)
6✔
1104
                length := uint64(0)
6✔
1105

6✔
1106
                query := req.URL.Query()
6✔
1107
                if offsetStr := query.Get("offset"); offsetStr != "" {
9✔
1108
                        o, err := strconv.ParseUint(offsetStr, 10, 64)
3✔
1109
                        if err != nil {
4✔
1110
                                http.Error(w, fmt.Sprintf("invalid offset %q: %v", offsetStr, err), http.StatusBadRequest)
1✔
1111
                                return
1✔
1112
                        }
1✔
1113
                        offset = o
2✔
1114
                }
1115
                if lengthStr := query.Get("length"); lengthStr != "" {
8✔
1116
                        l, err := strconv.ParseUint(lengthStr, 10, 64)
3✔
1117
                        if err != nil {
4✔
1118
                                http.Error(w, fmt.Sprintf("invalid length %q: %v", lengthStr, err), http.StatusBadRequest)
1✔
1119
                                return
1✔
1120
                        }
1✔
1121
                        length = l
2✔
1122
                }
1123

1124
                stream, err := client.Read(req.Context(), &nbdv1.ReadRequest{
4✔
1125
                        ExportName: exportName,
4✔
1126
                        Offset:     offset,
4✔
1127
                        Length:     length,
4✔
1128
                })
4✔
1129
                if err != nil {
5✔
1130
                        http.Error(w, fmt.Sprintf("Failed to call read for export: %s", exportName), http.StatusInternalServerError)
1✔
1131
                        return
1✔
1132
                }
1✔
1133

1134
                w.Header().Set("Content-Type", "application/octet-stream")
3✔
1135
                for {
8✔
1136
                        chunk, err := stream.Recv()
5✔
1137
                        if errors.Is(err, io.EOF) {
8✔
1138
                                break
3✔
1139
                        }
1140
                        if err != nil {
2✔
1141
                                log.Log.Reason(err).Error("Tunnel stream interrupted")
×
1142
                                panic(http.ErrAbortHandler)
×
1143
                        }
1144
                        if _, err := w.Write(chunk.Data); err != nil {
2✔
1145
                                log.Log.Reason(err).Error("HTTP client disconnected during stream")
×
1146
                                return
×
1147
                        }
×
1148
                        if f, ok := w.(http.Flusher); ok {
4✔
1149
                                f.Flush()
2✔
1150
                        }
2✔
1151
                }
1152
        })
1153
}
1154

1155
func collectMapPage(stream nbdv1.NBD_MapClient, pageSize int) ([]ExportMapExtent, *uint64, error) {
9✔
1156
        var extents []ExportMapExtent
9✔
1157
        for {
26✔
1158
                msg, err := stream.Recv()
17✔
1159
                if errors.Is(err, io.EOF) {
22✔
1160
                        return extents, nil, nil
5✔
1161
                }
5✔
1162
                if err != nil {
14✔
1163
                        return nil, nil, err
2✔
1164
                }
2✔
1165
                for _, e := range msg.Extents {
20✔
1166
                        if len(extents) >= pageSize {
12✔
1167
                                return extents, &e.Offset, nil
2✔
1168
                        }
2✔
1169
                        extents = append(extents, ExportMapExtent{
8✔
1170
                                Offset:      e.Offset,
8✔
1171
                                Length:      e.Length,
8✔
1172
                                Type:        e.Flags,
8✔
1173
                                Description: e.Description,
8✔
1174
                        })
8✔
1175
                }
1176
        }
1177
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc