Compare commits
2
Commits
5768734891
...
fcc9f34b83
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fcc9f34b83 | ||
|
|
8496404d8b |
@@ -22,6 +22,26 @@ the tag exists and is exercised; what is left is merging `next` to
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
- 2026-10-07: Made a backup notice a file rewritten with its size
|
||||||
|
unchanged and a new mtime in the same second as the one in the index
|
||||||
|
([issue #226](https://git.eeqj.de/sneak/vaultik/issues/226)). The
|
||||||
|
`files` table held mtime in whole seconds and the scanner compared
|
||||||
|
whole seconds, so every later snapshot kept the old content. A new
|
||||||
|
`mtime_nsec` column holds the nanoseconds within the second that
|
||||||
|
`mtime` holds, and the scanner compares the full mtime. A local index
|
||||||
|
created before the change lacks the column and is rebuilt with
|
||||||
|
`vaultik database delete` and a full backup.
|
||||||
|
|
||||||
|
- 2026-10-07: Made taking the process-wide lock atomic
|
||||||
|
([issue #227](https://git.eeqj.de/sneak/vaultik/issues/227)). The lock
|
||||||
|
read `vaultik.pid`, checked whether that PID was alive and then wrote
|
||||||
|
its own, so two writers started together could both pass the check and
|
||||||
|
both run. It is now an `flock` on `vaultik.pid`, held until the run
|
||||||
|
ends; the kernel drops it when the process exits, so a crash leaves no
|
||||||
|
lock behind. A clean exit now empties the file instead of deleting it,
|
||||||
|
because deleting it would let two later runs each lock a different
|
||||||
|
file.
|
||||||
|
|
||||||
- 2026-10-06: Made `snapshot remove --json` write only its document to
|
- 2026-10-06: Made `snapshot remove --json` write only its document to
|
||||||
stdout when the destination store cannot be reached
|
stdout when the destination store cannot be reached
|
||||||
([issue #251](https://git.eeqj.de/sneak/vaultik/issues/251)). Its
|
([issue #251](https://git.eeqj.de/sneak/vaultik/issues/251)). Its
|
||||||
|
|||||||
+2
-1
@@ -36,7 +36,8 @@ Stores metadata about files in the filesystem being backed up.
|
|||||||
**Columns:**
|
**Columns:**
|
||||||
- `id` (TEXT PRIMARY KEY) - UUID for the file record
|
- `id` (TEXT PRIMARY KEY) - UUID for the file record
|
||||||
- `path` (TEXT NOT NULL UNIQUE) - Absolute file path
|
- `path` (TEXT NOT NULL UNIQUE) - Absolute file path
|
||||||
- `mtime` (INTEGER NOT NULL) - Modification time as Unix timestamp
|
- `mtime` (INTEGER NOT NULL) - Modification time, whole seconds since the Unix epoch
|
||||||
|
- `mtime_nsec` (INTEGER NOT NULL) - Nanoseconds within that second, 0 to 999999999
|
||||||
- `size` (INTEGER NOT NULL) - File size in bytes
|
- `size` (INTEGER NOT NULL) - File size in bytes
|
||||||
- `mode` (INTEGER NOT NULL) - Unix file permissions and type
|
- `mode` (INTEGER NOT NULL) - Unix file permissions and type
|
||||||
- `uid` (INTEGER NOT NULL) - User ID of file owner
|
- `uid` (INTEGER NOT NULL) - User ID of file owner
|
||||||
|
|||||||
+28
-19
@@ -33,11 +33,13 @@ func (r *FileRepository) Create(ctx context.Context, tx *sql.Tx, file *File) err
|
|||||||
}
|
}
|
||||||
|
|
||||||
query := `
|
query := `
|
||||||
INSERT INTO files (id, path, source_path, mtime, size, mode, uid, gid, link_target)
|
INSERT INTO files
|
||||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
|
(id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
ON CONFLICT(path) DO UPDATE SET
|
ON CONFLICT(path) DO UPDATE SET
|
||||||
source_path = excluded.source_path,
|
source_path = excluded.source_path,
|
||||||
mtime = excluded.mtime,
|
mtime = excluded.mtime,
|
||||||
|
mtime_nsec = excluded.mtime_nsec,
|
||||||
size = excluded.size,
|
size = excluded.size,
|
||||||
mode = excluded.mode,
|
mode = excluded.mode,
|
||||||
uid = excluded.uid,
|
uid = excluded.uid,
|
||||||
@@ -54,16 +56,19 @@ func (r *FileRepository) Create(ctx context.Context, tx *sql.Tx, file *File) err
|
|||||||
if tx != nil {
|
if tx != nil {
|
||||||
LogSQL("Execute", query,
|
LogSQL("Execute", query,
|
||||||
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
||||||
file.MTime.Unix(), file.Size, file.Mode, file.UID, file.GID,
|
file.MTime.Unix(), file.MTime.Nanosecond(),
|
||||||
|
file.Size, file.Mode, file.UID, file.GID,
|
||||||
file.LinkTarget.String())
|
file.LinkTarget.String())
|
||||||
err = tx.QueryRowContext(ctx, query,
|
err = tx.QueryRowContext(ctx, query,
|
||||||
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
||||||
file.MTime.Unix(), file.Size, file.Mode, file.UID, file.GID,
|
file.MTime.Unix(), file.MTime.Nanosecond(),
|
||||||
|
file.Size, file.Mode, file.UID, file.GID,
|
||||||
file.LinkTarget.String()).Scan(&idStr)
|
file.LinkTarget.String()).Scan(&idStr)
|
||||||
} else {
|
} else {
|
||||||
err = r.db.QueryRowWithLog(ctx, query,
|
err = r.db.QueryRowWithLog(ctx, query,
|
||||||
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
file.ID.String(), file.Path.String(), file.SourcePath.String(),
|
||||||
file.MTime.Unix(), file.Size, file.Mode, file.UID, file.GID,
|
file.MTime.Unix(), file.MTime.Nanosecond(),
|
||||||
|
file.Size, file.Mode, file.UID, file.GID,
|
||||||
file.LinkTarget.String()).Scan(&idStr)
|
file.LinkTarget.String()).Scan(&idStr)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -84,7 +89,7 @@ func (r *FileRepository) Create(ctx context.Context, tx *sql.Tx, file *File) err
|
|||||||
// in the index.
|
// in the index.
|
||||||
func (r *FileRepository) GetByPath(ctx context.Context, path string) (*File, error) {
|
func (r *FileRepository) GetByPath(ctx context.Context, path string) (*File, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
WHERE path = ?
|
WHERE path = ?
|
||||||
`
|
`
|
||||||
@@ -104,7 +109,7 @@ func (r *FileRepository) GetByPath(ctx context.Context, path string) (*File, err
|
|||||||
// GetByID retrieves a file by its UUID
|
// GetByID retrieves a file by its UUID
|
||||||
func (r *FileRepository) GetByID(ctx context.Context, id types.FileID) (*File, error) {
|
func (r *FileRepository) GetByID(ctx context.Context, id types.FileID) (*File, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
WHERE id = ?
|
WHERE id = ?
|
||||||
`
|
`
|
||||||
@@ -127,7 +132,7 @@ func (r *FileRepository) GetByPathTx(
|
|||||||
ctx context.Context, tx *sql.Tx, path string,
|
ctx context.Context, tx *sql.Tx, path string,
|
||||||
) (*File, error) {
|
) (*File, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
WHERE path = ?
|
WHERE path = ?
|
||||||
`
|
`
|
||||||
@@ -158,13 +163,14 @@ func (r *FileRepository) ListModifiedSince(
|
|||||||
ctx context.Context, since time.Time,
|
ctx context.Context, since time.Time,
|
||||||
) ([]*File, error) {
|
) ([]*File, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
WHERE mtime >= ?
|
WHERE (mtime, mtime_nsec) >= (?, ?)
|
||||||
ORDER BY path
|
ORDER BY path
|
||||||
`
|
`
|
||||||
|
|
||||||
rows, err := r.db.conn.QueryContext(ctx, query, since.Unix())
|
rows, err := r.db.conn.QueryContext(ctx, query,
|
||||||
|
since.Unix(), since.Nanosecond())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("querying files: %w", err)
|
return nil, fmt.Errorf("querying files: %w", err)
|
||||||
}
|
}
|
||||||
@@ -239,7 +245,7 @@ func (r *FileRepository) ListUnderPath(
|
|||||||
|
|
||||||
// LIKE would ignore ASCII case and treat _ and % in path as wildcards.
|
// LIKE would ignore ASCII case and treat _ and % in path as wildcards.
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
WHERE path = ? OR substr(path, 1, length(?)) = ?
|
WHERE path = ? OR substr(path, 1, length(?)) = ?
|
||||||
ORDER BY path
|
ORDER BY path
|
||||||
@@ -324,7 +330,7 @@ func (r *FileRepository) ListIDsWithChunksNotInUploadedBlobs(
|
|||||||
// ListAll returns all files in the database
|
// ListAll returns all files in the database
|
||||||
func (r *FileRepository) ListAll(ctx context.Context) ([]*File, error) {
|
func (r *FileRepository) ListAll(ctx context.Context) ([]*File, error) {
|
||||||
query := `
|
query := `
|
||||||
SELECT id, path, source_path, mtime, size, mode, uid, gid, link_target
|
SELECT id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target
|
||||||
FROM files
|
FROM files
|
||||||
ORDER BY path
|
ORDER BY path
|
||||||
`
|
`
|
||||||
@@ -365,7 +371,7 @@ func (r *FileRepository) CreateBatch(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Each files row binds this many SQL variables.
|
// Each files row binds this many SQL variables.
|
||||||
const fileCols = 9
|
const fileCols = 10
|
||||||
|
|
||||||
// Batch at 100 rows to be safe with SQLite's variable limit.
|
// Batch at 100 rows to be safe with SQLite's variable limit.
|
||||||
const batchSize = 100
|
const batchSize = 100
|
||||||
@@ -376,7 +382,7 @@ func (r *FileRepository) CreateBatch(
|
|||||||
batch := files[i:end]
|
batch := files[i:end]
|
||||||
|
|
||||||
query := `INSERT INTO files
|
query := `INSERT INTO files
|
||||||
(id, path, source_path, mtime, size, mode, uid, gid, link_target)
|
(id, path, source_path, mtime, mtime_nsec, size, mode, uid, gid, link_target)
|
||||||
VALUES `
|
VALUES `
|
||||||
|
|
||||||
args := make([]any, 0, len(batch)*fileCols)
|
args := make([]any, 0, len(batch)*fileCols)
|
||||||
@@ -388,11 +394,12 @@ func (r *FileRepository) CreateBatch(
|
|||||||
querySb325.WriteString(", ")
|
querySb325.WriteString(", ")
|
||||||
}
|
}
|
||||||
|
|
||||||
querySb325.WriteString("(?, ?, ?, ?, ?, ?, ?, ?, ?)")
|
querySb325.WriteString("(?, ?, ?, ?, ?, ?, ?, ?, ?, ?)")
|
||||||
|
|
||||||
args = append(args,
|
args = append(args,
|
||||||
f.ID.String(), f.Path.String(), f.SourcePath.String(),
|
f.ID.String(), f.Path.String(), f.SourcePath.String(),
|
||||||
f.MTime.Unix(), f.Size, f.Mode, f.UID, f.GID,
|
f.MTime.Unix(), f.MTime.Nanosecond(),
|
||||||
|
f.Size, f.Mode, f.UID, f.GID,
|
||||||
f.LinkTarget.String())
|
f.LinkTarget.String())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -401,6 +408,7 @@ func (r *FileRepository) CreateBatch(
|
|||||||
query += ` ON CONFLICT(path) DO UPDATE SET
|
query += ` ON CONFLICT(path) DO UPDATE SET
|
||||||
source_path = excluded.source_path,
|
source_path = excluded.source_path,
|
||||||
mtime = excluded.mtime,
|
mtime = excluded.mtime,
|
||||||
|
mtime_nsec = excluded.mtime_nsec,
|
||||||
size = excluded.size,
|
size = excluded.size,
|
||||||
mode = excluded.mode,
|
mode = excluded.mode,
|
||||||
uid = excluded.uid,
|
uid = excluded.uid,
|
||||||
@@ -460,7 +468,7 @@ func (r *FileRepository) scanFileFrom(row fileRowScanner) (*File, error) {
|
|||||||
var (
|
var (
|
||||||
file File
|
file File
|
||||||
idStr, pathStr, sourcePathStr string
|
idStr, pathStr, sourcePathStr string
|
||||||
mtimeUnix int64
|
mtimeUnix, mtimeNsec int64
|
||||||
linkTarget sql.NullString
|
linkTarget sql.NullString
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -469,6 +477,7 @@ func (r *FileRepository) scanFileFrom(row fileRowScanner) (*File, error) {
|
|||||||
&pathStr,
|
&pathStr,
|
||||||
&sourcePathStr,
|
&sourcePathStr,
|
||||||
&mtimeUnix,
|
&mtimeUnix,
|
||||||
|
&mtimeNsec,
|
||||||
&file.Size,
|
&file.Size,
|
||||||
&file.Mode,
|
&file.Mode,
|
||||||
&file.UID,
|
&file.UID,
|
||||||
@@ -487,7 +496,7 @@ func (r *FileRepository) scanFileFrom(row fileRowScanner) (*File, error) {
|
|||||||
file.Path = types.FilePath(pathStr)
|
file.Path = types.FilePath(pathStr)
|
||||||
file.SourcePath = types.SourcePath(sourcePathStr)
|
file.SourcePath = types.SourcePath(sourcePathStr)
|
||||||
|
|
||||||
file.MTime = time.Unix(mtimeUnix, 0).UTC()
|
file.MTime = time.Unix(mtimeUnix, mtimeNsec).UTC()
|
||||||
if linkTarget.Valid {
|
if linkTarget.Valid {
|
||||||
file.LinkTarget = types.FilePath(linkTarget.String)
|
file.LinkTarget = types.FilePath(linkTarget.String)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -252,6 +252,115 @@ func TestFileRepositorySymlink(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// An mtime after 2262 or before 1678 does not fit in int64 nanoseconds
|
||||||
|
// since the epoch, and must still come back from the database unchanged.
|
||||||
|
func TestFileRepositoryMTimeOutsideInt64NanosecondRange(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
db, cleanup := setupTestDB(t)
|
||||||
|
defer cleanup()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
repo := database.NewFileRepository(db)
|
||||||
|
|
||||||
|
mtimes := []time.Time{
|
||||||
|
time.Date(2300, time.January, 1, 0, 0, 0, 123456789, time.UTC),
|
||||||
|
time.Date(1601, time.January, 1, 0, 0, 0, 987654321, time.UTC),
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, mtime := range mtimes {
|
||||||
|
created := &database.File{
|
||||||
|
Path: types.FilePath("/created-" + mtime.Format(time.RFC3339Nano)),
|
||||||
|
MTime: mtime,
|
||||||
|
}
|
||||||
|
|
||||||
|
err := repo.Create(ctx, nil, created)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to create file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
batched := &database.File{
|
||||||
|
ID: types.NewFileID(),
|
||||||
|
Path: types.FilePath("/batched-" + mtime.Format(time.RFC3339Nano)),
|
||||||
|
MTime: mtime,
|
||||||
|
}
|
||||||
|
|
||||||
|
err = repo.CreateBatch(ctx, nil, []*database.File{batched})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to batch create file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, path := range []types.FilePath{created.Path, batched.Path} {
|
||||||
|
retrieved, err := repo.GetByPath(ctx, path.String())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to get file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if !retrieved.MTime.Equal(mtime) {
|
||||||
|
t.Errorf("%s: mtime got %v, want %v",
|
||||||
|
path, retrieved.MTime, mtime)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A file already in the index and rewritten within the same second must get
|
||||||
|
// its new nanoseconds stored, through both Create and CreateBatch.
|
||||||
|
func TestFileRepositoryUpsertMTimeInSameSecond(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
db, cleanup := setupTestDB(t)
|
||||||
|
defer cleanup()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
repo := database.NewFileRepository(db)
|
||||||
|
|
||||||
|
indexed := time.Date(2026, time.October, 7, 12, 0, 0, 100000000, time.UTC)
|
||||||
|
rewritten := indexed.Add(800 * time.Millisecond)
|
||||||
|
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
upsert func(file *database.File) error
|
||||||
|
}{
|
||||||
|
{"Create", func(file *database.File) error {
|
||||||
|
return repo.Create(ctx, nil, file)
|
||||||
|
}},
|
||||||
|
{"CreateBatch", func(file *database.File) error {
|
||||||
|
return repo.CreateBatch(ctx, nil, []*database.File{file})
|
||||||
|
}},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
file := &database.File{
|
||||||
|
ID: types.NewFileID(),
|
||||||
|
Path: types.FilePath("/" + tt.name),
|
||||||
|
MTime: indexed,
|
||||||
|
}
|
||||||
|
|
||||||
|
err := tt.upsert(file)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("%s: failed to create file: %v", tt.name, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
file.MTime = rewritten
|
||||||
|
|
||||||
|
err = tt.upsert(file)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("%s: failed to update file: %v", tt.name, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
retrieved, err := repo.GetByPath(ctx, file.Path.String())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("%s: failed to get file: %v", tt.name, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if !retrieved.MTime.Equal(rewritten) {
|
||||||
|
t.Errorf("%s: mtime got %v, want %v",
|
||||||
|
tt.name, retrieved.MTime, rewritten)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestFileRepositoryTransaction(t *testing.T) {
|
func TestFileRepositoryTransaction(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -606,8 +606,7 @@ func TestTimezoneHandling(t *testing.T) {
|
|||||||
t.Skip("timezone not available")
|
t.Skip("timezone not available")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Use Truncate to remove sub-second precision since we store as Unix timestamps
|
nyTime := time.Now().In(loc)
|
||||||
nyTime := time.Now().In(loc).Truncate(time.Second)
|
|
||||||
file := &File{
|
file := &File{
|
||||||
Path: "/timezone-test.txt",
|
Path: "/timezone-test.txt",
|
||||||
MTime: nyTime,
|
MTime: nyTime,
|
||||||
|
|||||||
@@ -6,7 +6,8 @@ CREATE TABLE IF NOT EXISTS files (
|
|||||||
id TEXT PRIMARY KEY, -- UUID
|
id TEXT PRIMARY KEY, -- UUID
|
||||||
path TEXT NOT NULL UNIQUE,
|
path TEXT NOT NULL UNIQUE,
|
||||||
source_path TEXT NOT NULL DEFAULT '', -- The source directory this file came from (for restore path stripping)
|
source_path TEXT NOT NULL DEFAULT '', -- The source directory this file came from (for restore path stripping)
|
||||||
mtime INTEGER NOT NULL,
|
mtime INTEGER NOT NULL, -- whole seconds since the Unix epoch
|
||||||
|
mtime_nsec INTEGER NOT NULL, -- nanoseconds within that second, 0 to 999999999
|
||||||
size INTEGER NOT NULL,
|
size INTEGER NOT NULL,
|
||||||
mode INTEGER NOT NULL,
|
mode INTEGER NOT NULL,
|
||||||
uid INTEGER NOT NULL,
|
uid INTEGER NOT NULL,
|
||||||
|
|||||||
+67
-52
@@ -10,15 +10,18 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"syscall"
|
|
||||||
|
"golang.org/x/sys/unix"
|
||||||
)
|
)
|
||||||
|
|
||||||
// ErrAlreadyRunning indicates another vaultik instance is running.
|
// ErrAlreadyRunning indicates another vaultik instance is running.
|
||||||
var ErrAlreadyRunning = errors.New("another vaultik instance is already running")
|
var ErrAlreadyRunning = errors.New("another vaultik instance is already running")
|
||||||
|
|
||||||
// Lock represents an acquired PID lock.
|
// Lock represents an acquired PID lock: an flock(2) on the PID file,
|
||||||
|
// held while the file stays open. The kernel drops it when the process
|
||||||
|
// exits, however it exits, so a crashed run never leaves the lock held.
|
||||||
type Lock struct {
|
type Lock struct {
|
||||||
path string
|
file *os.File
|
||||||
}
|
}
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -29,10 +32,9 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// Acquire attempts to acquire a PID lock in the specified directory.
|
// Acquire attempts to acquire a PID lock in the specified directory.
|
||||||
// If the lock file exists and the process is still running, it returns
|
// If another process holds the lock, it returns ErrAlreadyRunning with
|
||||||
// ErrAlreadyRunning with details about the existing process.
|
// that process's PID. On success, it writes the current PID to the lock
|
||||||
// On success, it writes the current PID to the lock file and returns
|
// file and returns a Lock that must be released with Release().
|
||||||
// a Lock that must be released with Release().
|
|
||||||
func Acquire(lockDir string) (*Lock, error) {
|
func Acquire(lockDir string) (*Lock, error) {
|
||||||
// Ensure lock directory exists
|
// Ensure lock directory exists
|
||||||
err := os.MkdirAll(lockDir, lockDirPerm)
|
err := os.MkdirAll(lockDir, lockDirPerm)
|
||||||
@@ -42,56 +44,82 @@ func Acquire(lockDir string) (*Lock, error) {
|
|||||||
|
|
||||||
lockPath := filepath.Join(lockDir, "vaultik.pid")
|
lockPath := filepath.Join(lockDir, "vaultik.pid")
|
||||||
|
|
||||||
// Check for existing lock
|
// No O_TRUNC: the file may hold the PID of the process that has the
|
||||||
existingPID, err := readPIDFile(lockPath)
|
// lock, which the error below reports.
|
||||||
if err == nil {
|
file, err := os.OpenFile( //nolint:gosec // G304: path is our own lock file
|
||||||
// Lock file exists, check if process is running
|
lockPath, os.O_RDWR|os.O_CREATE, pidFilePerm)
|
||||||
if isProcessRunning(existingPID) {
|
|
||||||
return nil, fmt.Errorf("%w (PID %d)", ErrAlreadyRunning, existingPID)
|
|
||||||
}
|
|
||||||
// Process is not running, stale lock file - we can take over
|
|
||||||
}
|
|
||||||
|
|
||||||
// Write our PID
|
|
||||||
pid := os.Getpid()
|
|
||||||
|
|
||||||
err = os.WriteFile(lockPath, []byte(strconv.Itoa(pid)), pidFilePerm)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("writing PID file: %w", err)
|
return nil, fmt.Errorf("opening PID file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return &Lock{path: lockPath}, nil
|
err = unix.Flock(int(file.Fd()), unix.LOCK_EX|unix.LOCK_NB)
|
||||||
|
if err != nil {
|
||||||
|
_ = file.Close()
|
||||||
|
|
||||||
|
if errors.Is(err, unix.EWOULDBLOCK) {
|
||||||
|
return nil, alreadyRunningError(lockPath)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil, fmt.Errorf("locking PID file: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = writePID(file)
|
||||||
|
if err != nil {
|
||||||
|
_ = file.Close()
|
||||||
|
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return &Lock{file: file}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Release removes the PID lock file.
|
// Release empties the PID file and closes it, which drops the lock.
|
||||||
// It is safe to call Release multiple times.
|
// It is safe to call Release multiple times.
|
||||||
func (l *Lock) Release() error {
|
func (l *Lock) Release() error {
|
||||||
if l == nil || l.path == "" {
|
if l == nil || l.file == nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Verify we still own the lock (our PID is in the file)
|
file := l.file
|
||||||
existingPID, err := readPIDFile(l.path)
|
l.file = nil
|
||||||
|
|
||||||
|
// Do not remove the file here. A process that opened it a moment
|
||||||
|
// earlier could then lock the removed file while another creates and
|
||||||
|
// locks a new one, and both would run.
|
||||||
|
truncateErr := file.Truncate(0)
|
||||||
|
closeErr := file.Close()
|
||||||
|
|
||||||
|
return errors.Join(truncateErr, closeErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
// writePID replaces the contents of the locked PID file with the current
|
||||||
|
// PID.
|
||||||
|
func writePID(file *os.File) error {
|
||||||
|
err := file.Truncate(0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// File already gone or unreadable - that's fine
|
return fmt.Errorf("truncating PID file: %w", err)
|
||||||
return nil //nolint:nilerr // unreadable lock file means nothing to release
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if existingPID != os.Getpid() {
|
_, err = file.WriteAt([]byte(strconv.Itoa(os.Getpid())), 0)
|
||||||
// Someone else wrote to our lock file - don't remove it
|
if err != nil {
|
||||||
return nil
|
return fmt.Errorf("writing PID file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = os.Remove(l.path)
|
|
||||||
if err != nil && !os.IsNotExist(err) {
|
|
||||||
return fmt.Errorf("removing PID file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
l.path = "" // Prevent double-release
|
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// alreadyRunningError reports that another process holds the lock,
|
||||||
|
// naming its PID when the file holds one. The holder writes its PID just
|
||||||
|
// after it locks, so the file can briefly be empty.
|
||||||
|
func alreadyRunningError(lockPath string) error {
|
||||||
|
pid, err := readPIDFile(lockPath)
|
||||||
|
if err != nil {
|
||||||
|
return ErrAlreadyRunning
|
||||||
|
}
|
||||||
|
|
||||||
|
return fmt.Errorf("%w (PID %d)", ErrAlreadyRunning, pid)
|
||||||
|
}
|
||||||
|
|
||||||
// readPIDFile reads and parses the PID from a lock file.
|
// readPIDFile reads and parses the PID from a lock file.
|
||||||
func readPIDFile(path string) (int, error) {
|
func readPIDFile(path string) (int, error) {
|
||||||
data, err := os.ReadFile(path) //nolint:gosec // G304: path is our own lock file
|
data, err := os.ReadFile(path) //nolint:gosec // G304: path is our own lock file
|
||||||
@@ -106,16 +134,3 @@ func readPIDFile(path string) (int, error) {
|
|||||||
|
|
||||||
return pid, nil
|
return pid, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// isProcessRunning checks if a process with the given PID is running.
|
|
||||||
func isProcessRunning(pid int) bool {
|
|
||||||
process, err := os.FindProcess(pid)
|
|
||||||
if err != nil {
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
// On Unix, FindProcess always succeeds. We need to send signal 0 to check.
|
|
||||||
err = process.Signal(syscall.Signal(0))
|
|
||||||
|
|
||||||
return err == nil
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
@@ -33,9 +34,10 @@ func TestAcquireAndRelease(t *testing.T) {
|
|||||||
err = lock.Release()
|
err = lock.Release()
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
// Verify PID file is gone
|
// Verify PID file is empty
|
||||||
_, err = os.Stat(pidPath)
|
data, err = os.ReadFile(pidPath) //nolint:gosec // G304: test's own temp file
|
||||||
assert.True(t, os.IsNotExist(err))
|
require.NoError(t, err)
|
||||||
|
assert.Empty(t, data)
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestAcquireBlocksSecondInstance(t *testing.T) {
|
func TestAcquireBlocksSecondInstance(t *testing.T) {
|
||||||
@@ -55,6 +57,64 @@ func TestAcquireBlocksSecondInstance(t *testing.T) {
|
|||||||
lock2, err := pidlock.Acquire(tmpDir)
|
lock2, err := pidlock.Acquire(tmpDir)
|
||||||
require.ErrorIs(t, err, pidlock.ErrAlreadyRunning)
|
require.ErrorIs(t, err, pidlock.ErrAlreadyRunning)
|
||||||
assert.Nil(t, lock2)
|
assert.Nil(t, lock2)
|
||||||
|
|
||||||
|
// Once the first lock is released, the next Acquire succeeds
|
||||||
|
require.NoError(t, lock1.Release())
|
||||||
|
|
||||||
|
lock3, err := pidlock.Acquire(tmpDir)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NoError(t, lock3.Release())
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestConcurrentAcquireAdmitsOne starts many Acquire calls at the same
|
||||||
|
// moment, as two cron entries firing together would, and checks that
|
||||||
|
// exactly one of them gets the lock.
|
||||||
|
func TestConcurrentAcquireAdmitsOne(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
const callers = 50
|
||||||
|
|
||||||
|
tmpDir := t.TempDir()
|
||||||
|
start := make(chan struct{})
|
||||||
|
|
||||||
|
var (
|
||||||
|
mu sync.Mutex
|
||||||
|
acquired []*pidlock.Lock
|
||||||
|
failures []error
|
||||||
|
wg sync.WaitGroup
|
||||||
|
)
|
||||||
|
|
||||||
|
for range callers {
|
||||||
|
wg.Go(func() {
|
||||||
|
<-start
|
||||||
|
|
||||||
|
lock, err := pidlock.Acquire(tmpDir)
|
||||||
|
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
failures = append(failures, err)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
acquired = append(acquired, lock)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
close(start)
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
for _, lock := range acquired {
|
||||||
|
require.NoError(t, lock.Release())
|
||||||
|
}
|
||||||
|
|
||||||
|
assert.Len(t, acquired, 1, "exactly one caller should hold the lock")
|
||||||
|
|
||||||
|
for _, err := range failures {
|
||||||
|
require.ErrorIs(t, err, pidlock.ErrAlreadyRunning)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestAcquireWithStaleLock(t *testing.T) {
|
func TestAcquireWithStaleLock(t *testing.T) {
|
||||||
|
|||||||
@@ -1221,7 +1221,7 @@ func (s *Scanner) checkFileInMemory(
|
|||||||
|
|
||||||
// Check if file has changed
|
// Check if file has changed
|
||||||
if existingFile.Size != file.Size ||
|
if existingFile.Size != file.Size ||
|
||||||
existingFile.MTime.Unix() != file.MTime.Unix() ||
|
!existingFile.MTime.Equal(file.MTime) ||
|
||||||
existingFile.Mode != file.Mode ||
|
existingFile.Mode != file.Mode ||
|
||||||
existingFile.UID != file.UID ||
|
existingFile.UID != file.UID ||
|
||||||
existingFile.GID != file.GID {
|
existingFile.GID != file.GID {
|
||||||
|
|||||||
@@ -0,0 +1,73 @@
|
|||||||
|
package vaultik_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"path/filepath"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/spf13/afero"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"sneak.berlin/go/vaultik/internal/database"
|
||||||
|
"sneak.berlin/go/vaultik/internal/log"
|
||||||
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
|
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A file rewritten with its size unchanged and a new mtime in the same
|
||||||
|
// second as the mtime the index holds must still be backed up. See
|
||||||
|
// https://git.eeqj.de/sneak/vaultik/issues/226.
|
||||||
|
//
|
||||||
|
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||||
|
func TestBackupOfSameSecondRewriteRestoresNewContent(t *testing.T) {
|
||||||
|
log.Initialize(log.Config{})
|
||||||
|
|
||||||
|
fs := afero.NewOsFs()
|
||||||
|
tempDir := t.TempDir()
|
||||||
|
dataDir := filepath.Join(tempDir, "src")
|
||||||
|
storeDir := filepath.Join(tempDir, "remote")
|
||||||
|
restoreDir := filepath.Join(tempDir, "restored")
|
||||||
|
dbPath := filepath.Join(tempDir, "index.sqlite")
|
||||||
|
rewrittenPath := filepath.Join(dataDir, "small.txt")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
files := writeFaultSourceTree(t, fs, dataDir)
|
||||||
|
cfg := changedFileConfig(dataDir, dbPath)
|
||||||
|
|
||||||
|
firstMTime := time.Date(2026, time.January, 2, 3, 4, 5, 0, time.UTC).
|
||||||
|
Add(100 * time.Millisecond)
|
||||||
|
secondMTime := firstMTime.Add(800 * time.Millisecond)
|
||||||
|
|
||||||
|
require.NoError(t, fs.Chtimes(rewrittenPath, firstMTime, firstMTime))
|
||||||
|
|
||||||
|
store, err := storage.NewFileStorer(storeDir)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
db, err := database.New(ctx, dbPath)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
repos := database.NewRepositories(db)
|
||||||
|
v := newBackupVaultik(ctx, cfg, store, repos, db, fs)
|
||||||
|
|
||||||
|
require.NoError(t, backUp(v, "first"))
|
||||||
|
|
||||||
|
// Upper-casing ASCII text keeps its size.
|
||||||
|
files[rewrittenPath] = bytes.ToUpper(files[rewrittenPath])
|
||||||
|
require.NoError(t, afero.WriteFile(fs, rewrittenPath, files[rewrittenPath], 0o644))
|
||||||
|
require.NoError(t, fs.Chtimes(rewrittenPath, secondMTime, secondMTime))
|
||||||
|
|
||||||
|
require.NoError(t, backUp(v, "second"))
|
||||||
|
|
||||||
|
id := localSnapshotID(ctx, t, repos, "second")
|
||||||
|
require.NoError(t, db.Close())
|
||||||
|
|
||||||
|
reader := newReaderVaultik(ctx, cfg, store, nil, fs)
|
||||||
|
require.NoError(t, reader.Restore(&vaultik.RestoreOptions{
|
||||||
|
SnapshotID: id,
|
||||||
|
TargetDir: restoreDir,
|
||||||
|
Verify: true,
|
||||||
|
}))
|
||||||
|
|
||||||
|
assertRestoredTree(t, fs, restoreDir, files)
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user