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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,8 @@ Automatic mounting uses FUSE on Linux and WebDAV on macOS and Windows. macOS use

Mount commands start the companion runtime in the background, wait until the mount is ready, and then return a structured result. Use `ti fs unmount-file-system` or `ti fs-vault unmount-vault` to end a mount. The public CLI does not expose a foreground mount mode.

Mount readiness means the mounted path is readable: `ti fs mount-file-system` and `ti fs-vault mount-vault` keep polling the mount path until directory listing succeeds, bounded by `--ready-timeout` (default `30s`). If the mount path never becomes readable in time, the command fails with `fs.mount_ready_timeout`, the background mount is left running, and the mount locator is preserved so the matching unmount command can stop it.

Filesystem layers can fork copy-on-write child timelines without copying a workspace. Layer and checkpoint mounts require FUSE; checkpoint mounts are always read-only. Drive9 does not support combining recursive copy with a layer, so seed a directory tree through a writable layer mount, drain it, and then create the checkpoint:

```shell
Expand Down
4 changes: 4 additions & 0 deletions e2e/cli_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -825,6 +825,7 @@ func TestFSRemoteInventoryAndIDCredentialSelectionAcrossCommandFamilies(t *testi
baseEnv := []string{
"HOME=" + home,
"TI_DRIVE9_BIN=" + companion,
"TI_TEST_FAKE_MOUNT_READY=1",
"FAKE_DRIVE9_RECORD=" + recordPath,
"TI_ALLOW_TEST_ENDPOINTS=1",
"TI_TEST_FS_MANIFEST_URL=" + manifestServer.URL,
Expand Down Expand Up @@ -1097,6 +1098,7 @@ func TestFSConfigurationFreeAccess(t *testing.T) {
"HOME=" + home,
"TI_LOGGING=on",
"TI_DRIVE9_BIN=" + companion,
"TI_TEST_FAKE_MOUNT_READY=1",
"FAKE_DRIVE9_RECORD=" + recordPath,
"TI_ALLOW_TEST_ENDPOINTS=1",
"TI_TEST_FS_MANIFEST_URL=" + manifestServer.URL,
Expand Down Expand Up @@ -1207,6 +1209,7 @@ func TestFSLayerForkWorkflowCommands(t *testing.T) {
env := []string{
"HOME=" + home,
"TI_DRIVE9_BIN=" + companion,
"TI_TEST_FAKE_MOUNT_READY=1",
"FAKE_DRIVE9_RECORD=" + recordPath,
"TI_ALLOW_TEST_ENDPOINTS=1",
"TI_TEST_FS_MANIFEST_URL=" + manifestServer.URL,
Expand Down Expand Up @@ -1295,6 +1298,7 @@ func TestFSImportFileSystemToken(t *testing.T) {
baseEnv := []string{
"HOME=" + home,
"TI_DRIVE9_BIN=" + companion,
"TI_TEST_FAKE_MOUNT_READY=1",
"FAKE_DRIVE9_RECORD=" + recordPath,
"FAKE_DRIVE9_EXPECT_API_KEY=" + token,
"TI_ALLOW_TEST_ENDPOINTS=1",
Expand Down
4 changes: 4 additions & 0 deletions e2e/testdata/fake-drive9.go
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,10 @@ func main() {
}
if hasPrefix(args, "mount") {
fmt.Fprintln(os.Stderr, "drive9: mount mode: "+mountMode(args))
if mountPath := args[len(args)-1]; len(args) >= 2 && len(mountPath) > 0 && mountPath[0] != '-' {
_ = os.MkdirAll(mountPath, 0o755)
_ = os.WriteFile(mountPath+string(os.PathSeparator)+".drive9-mounted", []byte("fake mount ready\n"), 0o644)
}
return
}
}
Expand Down
6 changes: 6 additions & 0 deletions internal/fs/control.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ type Service struct {
Timeout time.Duration
FSReadyWaitTimeout time.Duration
FSReadyWaitPollInterval time.Duration
MountReadyPollInterval time.Duration
Debug bool
DebugWriter io.Writer
HomeDir string
Expand All @@ -37,6 +38,11 @@ type Service struct {
Stdin io.Reader
Stdout io.Writer
Stderr io.Writer

// mountPointActive overrides active-mount evidence detection. It exists
// for tests of unrelated mount behaviors whose fake companion cannot
// create a real kernel mount; production leaves it nil.
mountPointActive func(string) (bool, error)
}

type CreateFileSystemOptions struct {
Expand Down
6 changes: 6 additions & 0 deletions internal/fs/drive9_companion.go
Original file line number Diff line number Diff line change
Expand Up @@ -1193,6 +1193,9 @@ func (s Service) drive9MountVault(ctx context.Context, opts VaultMountOptions) (
_, _ = s.drive9Run(ctx, opts.Profile, []string{"umount", opts.MountPath}, false)
return MountResult{}, err
}
if err := s.waitForMountReady(ctx, opts.MountPath, opts.ReadyTimeout, fmt.Sprintf("; to stop it run: ti fs-vault unmount-vault --mount-path %q", opts.MountPath)); err != nil {
return MountResult{}, err
}
return MountResult{Status: "mounted", Profile: profileName(opts.Profile), FileSystemName: "vault", MountPath: opts.MountPath, RemotePath: "/n/vault", Driver: "fuse"}, nil
}

Expand Down Expand Up @@ -1296,6 +1299,9 @@ func (s Service) drive9MountFileSystem(ctx context.Context, opts MountFileSystem
_, _ = s.drive9Run(ctx, opts.Profile, []string{"umount", opts.MountPath}, false)
return MountResult{}, err
}
if err := s.waitForMountReady(ctx, opts.MountPath, opts.ReadyTimeout, fmt.Sprintf("; to stop it run: ti fs unmount-file-system --mount-path %q", opts.MountPath)); err != nil {
return MountResult{}, err
}
endpoint, _ := s.resolveFS(opts.Profile)
driver := drive9MountedDriver(result.Stderr, opts.Driver)
if driver == "" {
Expand Down
12 changes: 11 additions & 1 deletion internal/fs/drive9_companion_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -414,7 +414,9 @@ func TestDrive9LayerCheckpointMountIsReadOnlyAndRecorded(t *testing.T) {
companion, recordPath := buildFakeDrive9(t)
t.Setenv("TI_FAKE_DRIVE9_RECORD", recordPath)
mountPath := filepath.Join(t.TempDir(), "checkpoint")
result, err := testCompanionService(home, companion).MountFileSystem(context.Background(), MountFileSystemOptions{
checkpointService := testCompanionService(home, companion)
checkpointService.mountPointActive = trueMountEvidence
result, err := checkpointService.MountFileSystem(context.Background(), MountFileSystemOptions{
Profile: dataProfile(), MountPath: mountPath, RemotePath: "/research/q3-market", Driver: "fuse", LayerRef: "style-analyst", CheckpointID: "v5",
})
if err != nil {
Expand Down Expand Up @@ -535,6 +537,7 @@ func TestDrive9MountLocatorRoutesDrainAndUnmountWithoutCredentials(t *testing.T)
companion, recordPath := buildFakeDrive9(t)
t.Setenv("TI_FAKE_DRIVE9_RECORD", recordPath)
service := testCompanionService(home, companion)
service.mountPointActive = trueMountEvidence
mountPath := filepath.Join(t.TempDir(), "workspace")
profile := dataProfile()

Expand Down Expand Up @@ -599,6 +602,7 @@ func TestDrive9VaultMountUsesBackgroundMode(t *testing.T) {
companion, recordPath := buildFakeDrive9(t)
t.Setenv("TI_FAKE_DRIVE9_RECORD", recordPath)
service := testCompanionService(home, companion)
service.mountPointActive = trueMountEvidence
service.Stdout = &bytes.Buffer{}
service.Stderr = &bytes.Buffer{}
mountPath := filepath.Join(t.TempDir(), "vault")
Expand Down Expand Up @@ -631,6 +635,7 @@ func TestDrive9MountSuppressesCompanionSuccessChatter(t *testing.T) {
var stdout bytes.Buffer
var stderr bytes.Buffer
service := testCompanionService(home, companion)
service.mountPointActive = trueMountEvidence
service.Stdout = &stdout
service.Stderr = &stderr

Expand Down Expand Up @@ -673,6 +678,7 @@ func TestDrive9FailedUnmountPreservesMountLocator(t *testing.T) {
companion, recordPath := buildFakeDrive9(t)
t.Setenv("TI_FAKE_DRIVE9_RECORD", recordPath)
service := testCompanionService(home, companion)
service.mountPointActive = trueMountEvidence
mountPath := filepath.Join(t.TempDir(), "workspace")
if _, err := service.MountFileSystem(context.Background(), MountFileSystemOptions{
Profile: dataProfile(),
Expand Down Expand Up @@ -866,6 +872,10 @@ func main() {
fmt.Fprintln(os.Stderr, "mount: drive9 mount: background mount exited before becoming ready")
os.Exit(1)
}
if mountPath := args[len(args)-1]; len(args) >= 2 && len(mountPath) > 0 && mountPath[0] != '-' {
_ = os.MkdirAll(mountPath, 0o755)
_ = os.WriteFile(mountPath+string(os.PathSeparator)+".drive9-mounted", []byte("fake mount ready\n"), 0o644)
}
fmt.Fprintln(os.Stderr, "drive9: mount running in background")
fmt.Fprintln(os.Stderr, "drive9: unmount with drive9 umount /workspace")
case len(args) >= 3 && args[0] == "admin" && args[1] == "tenant" && args[2] == "delete":
Expand Down
149 changes: 149 additions & 0 deletions internal/fs/mountready.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
package fs

import (
"context"
"errors"
"fmt"
"io"
"os"
"time"

"github.com/tidbcloud/ti-cli/internal/apperr"
)

const (
defaultMountReadyTimeout = 30 * time.Second
defaultMountReadyPollInterval = 100 * time.Millisecond
// mountReadyProbeTimeout bounds a single probe so a wedged mount cannot
// block one probe forever. It is an upper bound only: every probe is also
// capped by the remaining --ready-timeout budget and the command context,
// so a short --ready-timeout never waits a full probe bound. It must stay
// well above a healthy cold first readdir: a WebDAV mount's initial
// directory listing crosses the companion proxy and the remote region and
// can legitimately take seconds.
mountReadyProbeTimeout = 10 * time.Second
)

// errMountReadyProbeExhausted marks a probe whose budget elapsed before the
// probe answered. It means "not ready yet", not a mount failure.
var errMountReadyProbeExhausted = errors.New("mount readiness probe budget exhausted")

// probeMountPointReady reports whether mountPath is an active mount that the
// kernel can list. A background mount is only usable once the mount exists,
// is visible in the mount table, and readdir is served from the mount root.
func probeMountPointReady(mountPath string, mounted func(string) (bool, error)) error {
info, err := os.Stat(mountPath)
if err != nil {
return err
}
if !info.IsDir() {
return fmt.Errorf("mount path %q is not a directory", mountPath)
}
active, err := mounted(mountPath)
if err != nil {
return fmt.Errorf("mount evidence for %q: %w", mountPath, err)
}
if !active {
return fmt.Errorf("mount path %q is not an active mount", mountPath)
}
dir, err := os.Open(mountPath)
if err != nil {
return err
}
defer dir.Close()
if _, err := dir.Readdirnames(1); err != nil && !errors.Is(err, io.EOF) {
return err
}
return nil
}

// probeMountPointOnce runs one probe bounded by budget and ctx. A blocked or
// wedged probe is abandoned when either expires; the abandoned goroutine
// finishes and closes its handle once the underlying syscall returns.
func probeMountPointOnce(ctx context.Context, mountPath string, budget time.Duration, mounted func(string) (bool, error)) error {
done := make(chan error, 1)
go func() {
done <- probeMountPointReady(mountPath, mounted)
}()
timer := time.NewTimer(budget)
defer timer.Stop()
select {
case err := <-done:
return err
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
return errMountReadyProbeExhausted
}
}

func (s Service) mountEvidence() func(string) (bool, error) {
if s.mountPointActive != nil {
return s.mountPointActive
}
return defaultMountPointActive
}

func (s Service) waitForMountReady(ctx context.Context, mountPath string, timeout time.Duration, stopHint string) error {
if timeout <= 0 {
timeout = defaultMountReadyTimeout
}
interval := s.MountReadyPollInterval
if interval <= 0 {
interval = defaultMountReadyPollInterval
}
mounted := s.mountEvidence()
deadline := time.Now().Add(timeout)
var lastErr error
for {
remaining := time.Until(deadline)
if remaining <= 0 {
break
}
budget := mountReadyProbeTimeout
if remaining < budget {
budget = remaining
}
err := probeMountPointOnce(ctx, mountPath, budget, mounted)
switch {
case err == nil:
return nil
case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded):
return mountReadyCanceled(mountPath, stopHint, err)
default:
lastErr = err
}
remaining = time.Until(deadline)
if remaining <= 0 {
break
}
wait := interval
if remaining < wait {
wait = remaining
}
timer := time.NewTimer(wait)
select {
case <-ctx.Done():
timer.Stop()
return mountReadyCanceled(mountPath, stopHint, ctx.Err())
case <-timer.C:
}
}
return apperr.Wrap(
"fs.mount_ready_timeout",
"runtime",
1,
fmt.Sprintf("background mount at %q did not become readable within %s; the mount is still running%s", mountPath, timeout, stopHint),
lastErr,
)
}

func mountReadyCanceled(mountPath, stopHint string, cause error) error {
return apperr.Wrap(
"fs.mount_ready_canceled",
"runtime",
1,
fmt.Sprintf("waiting for the background mount at %q to become readable was canceled; the mount is still running%s", mountPath, stopHint),
cause,
)
}
31 changes: 31 additions & 0 deletions internal/fs/mountready_darwin.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
//go:build darwin

package fs

import (
"os"
"path/filepath"
"syscall"
)

// defaultMountPointActive reports whether mountPath is the root of an active
// mount. On macOS a mount root sits on a different filesystem than its parent
// directory, so unequal st_dev values prove an active mount (FUSE, WebDAV,
// disk images); a plain subdirectory shares the parent's device.
func defaultMountPointActive(mountPath string) (bool, error) {
if testFakeMountReady() {
return true, nil
}
if mountPath == string(os.PathSeparator) {
return true, nil
}
info, err := os.Stat(mountPath)
if err != nil {
return false, err
}
parent, err := os.Stat(filepath.Dir(mountPath))
if err != nil {
return false, err
}
return info.Sys().(*syscall.Stat_t).Dev != parent.Sys().(*syscall.Stat_t).Dev, nil
}
70 changes: 70 additions & 0 deletions internal/fs/mountready_linux.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
//go:build linux

package fs

import (
"os"
"path/filepath"
"strconv"
"strings"
)

// defaultMountPointActive reports whether mountPath is the root of an active
// mount by checking /proc/self/mountinfo, the kernel's authoritative mount
// table. The companion mounts FUSE filesystems, which always appear there.
func defaultMountPointActive(mountPath string) (bool, error) {
if testFakeMountReady() {
return true, nil
}
candidates := mountPathCandidates(mountPath)
data, err := os.ReadFile("/proc/self/mountinfo")
if err != nil {
return false, err
}
for _, line := range strings.Split(string(data), "\n") {
fields := strings.Split(line, " ")
if len(fields) < 5 {
continue
}
if _, ok := candidates[unescapeMountinfoPath(fields[4])]; ok {
return true, nil
}
}
return false, nil
}

func mountPathCandidates(mountPath string) map[string]struct{} {
candidates := map[string]struct{}{}
for _, path := range []string{mountPath, filepath.Clean(mountPath)} {
if path != "" {
candidates[path] = struct{}{}
}
}
if resolved, err := filepath.EvalSymlinks(mountPath); err == nil {
candidates[resolved] = struct{}{}
}
if abs, err := filepath.Abs(mountPath); err == nil {
candidates[abs] = struct{}{}
}
return candidates
}

// unescapeMountinfoPath decodes the octal escapes (for example \040 for a
// space) that /proc/self/mountinfo uses for special characters.
func unescapeMountinfoPath(path string) string {
if !strings.Contains(path, "\\") {
return path
}
var out strings.Builder
for i := 0; i < len(path); i++ {
if path[i] == '\\' && i+3 < len(path) {
if value, err := strconv.ParseUint(path[i+1:i+4], 8, 8); err == nil {
out.WriteByte(byte(value))
i += 3
continue
}
}
out.WriteByte(path[i])
}
return out.String()
}
Loading
Loading