Files
unc-file-watcher/watcher.go
T

628 lines
15 KiB
Go

package main
import (
"archive/tar"
"compress/gzip"
"crypto/sha256"
"fmt"
"io"
"os"
"path"
"path/filepath"
"sort"
"strings"
"time"
)
type contentRecord struct {
path string
kind byte
size int64
sum [32]byte
}
func ensureExistingDir(path string) error {
info, err := os.Stat(path)
if err != nil {
return err
}
if !info.IsDir() {
return fmt.Errorf("%q 不是目录", path)
}
return nil
}
func ensureDir(path string) error {
info, err := os.Stat(path)
if err == nil {
if !info.IsDir() {
return fmt.Errorf("%q 不是目录", path)
}
return nil
}
if !os.IsNotExist(err) {
return err
}
return os.MkdirAll(path, 0o755)
}
func scanRule(
ruleName string,
root string,
recursive bool,
matches func(string) bool,
cfg config,
state map[string]fileState,
process processFunc,
logf func(string, ...any),
) error {
entries, err := collectCandidates(root, recursive, matches)
if err != nil {
return err
}
now := time.Now()
seen := make(map[string]struct{}, len(entries))
for _, entry := range entries {
sourcePath := entry.path
seen[sourcePath] = struct{}{}
snap := fileSnapshot{
size: entry.info.Size(),
modTime: entry.info.ModTime(),
}
key := snapshotKey(snap)
current := state[sourcePath]
if current.lastObserved != snap {
state[sourcePath] = fileState{
lastObserved: snap,
stableCount: 1,
}
logf("[%s] 发现新文件或文件有变化: %s", ruleName, sourcePath)
continue
}
if current.handledKey == key {
continue
}
current.stableCount++
current.lastObserved = snap
state[sourcePath] = current
if !isStableEnough(cfg, current.stableCount, snap, now) {
continue
}
decision, handled, err := process(sourcePath)
if err != nil {
logf("[%s] 处理失败: %s, err=%v", ruleName, sourcePath, err)
continue
}
if handled {
current.handledKey = key
} else {
current.handledKey = ""
}
if decision != current.lastDecision {
logf("[%s] %s: %s", ruleName, decision, sourcePath)
}
current.lastDecision = decision
state[sourcePath] = current
}
for path := range state {
if _, ok := seen[path]; !ok {
delete(state, path)
}
}
return nil
}
func collectCandidates(root string, recursive bool, matches func(string) bool) ([]candidate, error) {
if recursive {
var items []candidate
err := filepath.Walk(root, func(path string, info os.FileInfo, walkErr error) error {
if walkErr != nil {
return walkErr
}
if info.IsDir() {
return nil
}
if !matches(info.Name()) {
return nil
}
items = append(items, candidate{path: path, info: info})
return nil
})
return items, err
}
entries, err := os.ReadDir(root)
if err != nil {
return nil, err
}
items := make([]candidate, 0, len(entries))
for _, entry := range entries {
if entry.IsDir() {
continue
}
name := entry.Name()
if !matches(name) {
continue
}
info, err := entry.Info()
if err != nil {
return nil, err
}
items = append(items, candidate{
path: filepath.Join(root, name),
info: info,
})
}
return items, nil
}
func processArchiveCandidate(cfg config, sourcePath string) (string, bool, error) {
fileName := filepath.Base(sourcePath)
zipTarget := filepath.Join(cfg.archive.zipDir, fileName)
if pathExists(zipTarget) {
if cfg.mode == "move" {
same, err := filesHaveSameHash(sourcePath, zipTarget)
if err != nil {
return "", false, err
}
if same {
if err := os.Remove(sourcePath); err != nil {
return "", false, fmt.Errorf("remove duplicate archive source: %w", err)
}
return "已跳过传输,zip 目录已有同名压缩包且哈希一致,源文件已删除", true, nil
}
return "已跳过,zip 目录中存在同名压缩包但哈希不同", false, nil
}
return "已跳过,zip 目录中存在同名压缩包", false, nil
}
extractedDir := extractedDirName(fileName)
unzipTarget := filepath.Join(cfg.archive.unzipDir, extractedDir)
if isDir(unzipTarget) {
if cfg.mode == "move" {
same, err := archiveMatchesDirectory(sourcePath, unzipTarget, extractedDir)
if err != nil {
return "", false, err
}
if same {
if err := os.Remove(sourcePath); err != nil {
return "", false, fmt.Errorf("remove duplicate archive source: %w", err)
}
return "已跳过传输,unzip 目录已有同名解压目录且内容一致,源文件已删除", true, nil
}
return "已跳过,unzip 目录中存在同名解压目录但内容不同", false, nil
}
return "已跳过,unzip 目录中存在同名解压目录", false, nil
}
if err := transferFile(sourcePath, zipTarget, cfg.mode); err != nil {
return "", false, err
}
if cfg.mode == "move" {
return "已剪切到 zip 目录", true, nil
}
return "已复制到 zip 目录", true, nil
}
func processGenericCandidate(cfg config, sourcePath string) (string, bool, error) {
fileName := filepath.Base(sourcePath)
targetPath := filepath.Join(cfg.generic.targetDir, fileName)
if pathExists(targetPath) {
if cfg.mode == "move" {
return handleGenericDuplicateMove(sourcePath, targetPath, cfg.generic.dedupeMode)
}
decision, duplicate, err := evaluateGenericDuplicate(sourcePath, targetPath, cfg.generic.dedupeMode)
if err != nil {
return "", false, err
}
if duplicate {
return decision, false, nil
}
return decision, false, nil
}
if err := transferFile(sourcePath, targetPath, cfg.mode); err != nil {
return "", false, err
}
if cfg.mode == "move" {
return "已剪切到普通文件目标目录", true, nil
}
return "已复制到普通文件目标目录", true, nil
}
func evaluateGenericDuplicate(sourcePath, targetPath, dedupeMode string) (string, bool, error) {
decision, duplicate, _, err := analyzeGenericDuplicate(sourcePath, targetPath, dedupeMode)
return decision, duplicate, err
}
func analyzeGenericDuplicate(sourcePath, targetPath, dedupeMode string) (string, bool, bool, error) {
sourceInfo, err := os.Stat(sourcePath)
if err != nil {
return "", false, false, fmt.Errorf("stat source: %w", err)
}
targetInfo, err := os.Stat(targetPath)
if err != nil {
return "", false, false, fmt.Errorf("stat target: %w", err)
}
switch dedupeMode {
case genericDedupeNameHash:
if sourceInfo.Size() != targetInfo.Size() {
return "已跳过,同名文件已存在但大小不同", false, false, nil
}
same, err := filesHaveSameHash(sourcePath, targetPath)
if err != nil {
return "", false, false, err
}
if same {
return "已跳过,同名文件已存在且哈希相同", true, true, nil
}
return "已跳过,同名文件已存在但哈希不同", false, true, nil
case genericDedupeNameSize:
fallthrough
default:
if sourceInfo.Size() == targetInfo.Size() {
return "已跳过,同名文件已存在且大小相同", true, false, nil
}
return "已跳过,同名文件已存在但大小不同", false, false, nil
}
}
func handleGenericDuplicateMove(sourcePath, targetPath, dedupeMode string) (string, bool, error) {
decision, duplicate, hashVerified, err := analyzeGenericDuplicate(sourcePath, targetPath, dedupeMode)
if err != nil {
return "", false, err
}
if !duplicate {
return decision, false, nil
}
if !hashVerified {
same, err := filesHaveSameHash(sourcePath, targetPath)
if err != nil {
return "", false, err
}
if !same {
return "已跳过,同名文件大小相同但哈希不同", false, nil
}
}
if err := os.Remove(sourcePath); err != nil {
return "", false, fmt.Errorf("remove duplicate generic source: %w", err)
}
return "已跳过传输,目标目录已有同名文件且哈希一致,源文件已删除", true, nil
}
func filesHaveSameHash(sourcePath, targetPath string) (bool, error) {
sourceHash, err := fileSHA256(sourcePath)
if err != nil {
return false, fmt.Errorf("hash source: %w", err)
}
targetHash, err := fileSHA256(targetPath)
if err != nil {
return false, fmt.Errorf("hash target: %w", err)
}
return sourceHash == targetHash, nil
}
func archiveMatchesDirectory(sourcePath, targetPath, archiveRoot string) (bool, error) {
sourceRecords, err := archiveContentRecords(sourcePath, archiveRoot)
if err != nil {
return false, err
}
targetRecords, err := directoryContentRecords(targetPath)
if err != nil {
return false, err
}
if len(sourceRecords) != len(targetRecords) {
return false, nil
}
for i := range sourceRecords {
if sourceRecords[i] != targetRecords[i] {
return false, nil
}
}
return true, nil
}
func archiveContentRecords(archivePath, archiveRoot string) ([]contentRecord, error) {
file, err := os.Open(archivePath)
if err != nil {
return nil, fmt.Errorf("open archive: %w", err)
}
defer file.Close()
gzipReader, err := gzip.NewReader(file)
if err != nil {
return nil, fmt.Errorf("open gzip reader: %w", err)
}
defer gzipReader.Close()
records := make(map[string]contentRecord)
tarReader := tar.NewReader(gzipReader)
for {
header, err := tarReader.Next()
if err == io.EOF {
break
}
if err != nil {
return nil, fmt.Errorf("read archive entry: %w", err)
}
name := normalizeArchiveRecordPath(header.Name, archiveRoot)
switch header.Typeflag {
case tar.TypeReg, tar.TypeRegA, tar.TypeGNUSparse:
if name == "" {
return nil, fmt.Errorf("archive entry %q resolves to empty path", header.Name)
}
sum, err := readerSHA256(tarReader)
if err != nil {
return nil, fmt.Errorf("hash archive entry %q: %w", header.Name, err)
}
records[name] = contentRecord{
path: name,
kind: 'F',
size: header.Size,
sum: sum,
}
case tar.TypeSymlink:
if name == "" {
return nil, fmt.Errorf("archive symlink %q resolves to empty path", header.Name)
}
records[name] = contentRecord{
path: name,
kind: 'L',
sum: sha256.Sum256([]byte(filepath.ToSlash(header.Linkname))),
}
case tar.TypeDir, tar.TypeXHeader, tar.TypeXGlobalHeader, tar.TypeGNULongName, tar.TypeGNULongLink:
continue
default:
return nil, fmt.Errorf("unsupported archive entry type for %q: %d", header.Name, header.Typeflag)
}
}
return sortedContentRecords(records), nil
}
func directoryContentRecords(root string) ([]contentRecord, error) {
records := make(map[string]contentRecord)
err := filepath.Walk(root, func(currentPath string, info os.FileInfo, walkErr error) error {
if walkErr != nil {
return walkErr
}
if currentPath == root {
return nil
}
relativePath, err := filepath.Rel(root, currentPath)
if err != nil {
return fmt.Errorf("build relative path for %q: %w", currentPath, err)
}
relativePath = filepath.ToSlash(relativePath)
if info.Mode()&os.ModeSymlink != 0 {
linkTarget, err := os.Readlink(currentPath)
if err != nil {
return fmt.Errorf("read symlink %q: %w", currentPath, err)
}
records[relativePath] = contentRecord{
path: relativePath,
kind: 'L',
sum: sha256.Sum256([]byte(filepath.ToSlash(linkTarget))),
}
return nil
}
if info.IsDir() {
return nil
}
if !info.Mode().IsRegular() {
return fmt.Errorf("unsupported directory entry %q", currentPath)
}
sum, err := fileSHA256(currentPath)
if err != nil {
return fmt.Errorf("hash directory file %q: %w", currentPath, err)
}
records[relativePath] = contentRecord{
path: relativePath,
kind: 'F',
size: info.Size(),
sum: sum,
}
return nil
})
if err != nil {
return nil, err
}
return sortedContentRecords(records), nil
}
func normalizeArchiveRecordPath(name, archiveRoot string) string {
cleanPath := strings.TrimSpace(filepath.ToSlash(name))
cleanPath = strings.TrimPrefix(cleanPath, "./")
cleanPath = strings.TrimPrefix(cleanPath, "/")
cleanPath = path.Clean(cleanPath)
if cleanPath == "." {
return ""
}
normalizedRoot := path.Clean(strings.Trim(filepath.ToSlash(archiveRoot), "/"))
if normalizedRoot != "." && normalizedRoot != "" {
if cleanPath == normalizedRoot {
return ""
}
prefix := normalizedRoot + "/"
if strings.HasPrefix(cleanPath, prefix) {
cleanPath = strings.TrimPrefix(cleanPath, prefix)
}
}
if cleanPath == "." {
return ""
}
return cleanPath
}
func sortedContentRecords(records map[string]contentRecord) []contentRecord {
items := make([]contentRecord, 0, len(records))
for _, record := range records {
items = append(items, record)
}
sort.Slice(items, func(i, j int) bool {
if items[i].path == items[j].path {
return items[i].kind < items[j].kind
}
return items[i].path < items[j].path
})
return items
}
func fileSHA256(path string) ([32]byte, error) {
file, err := os.Open(path)
if err != nil {
return [32]byte{}, err
}
defer file.Close()
return readerSHA256(file)
}
func readerSHA256(reader io.Reader) ([32]byte, error) {
hasher := sha256.New()
if _, err := io.Copy(hasher, reader); err != nil {
return [32]byte{}, err
}
sumBytes := hasher.Sum(nil)
var sum [32]byte
copy(sum[:], sumBytes)
return sum, nil
}
func transferFile(sourcePath, targetPath, mode string) error {
if pathExists(targetPath) {
return fmt.Errorf("target already exists: %s", targetPath)
}
tempPath := targetPath + ".part"
if pathExists(tempPath) {
if err := os.Remove(tempPath); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove stale temp file: %w", err)
}
}
source, err := os.Open(sourcePath)
if err != nil {
return fmt.Errorf("open source: %w", err)
}
target, err := os.OpenFile(tempPath, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o644)
if err != nil {
_ = source.Close()
return fmt.Errorf("create temp target: %w", err)
}
if err := copyFileContents(target, source); err != nil {
_ = target.Close()
_ = source.Close()
_ = os.Remove(tempPath)
return err
}
if err := target.Close(); err != nil {
_ = source.Close()
_ = os.Remove(tempPath)
return fmt.Errorf("close temp target: %w", err)
}
if err := source.Close(); err != nil {
_ = os.Remove(tempPath)
return fmt.Errorf("close source: %w", err)
}
if err := os.Rename(tempPath, targetPath); err != nil {
_ = os.Remove(tempPath)
return fmt.Errorf("rename temp target: %w", err)
}
if mode == "move" {
if err := os.Remove(sourcePath); err != nil {
return fmt.Errorf("remove source after move: %w", err)
}
}
return nil
}
func copyFileContents(target *os.File, source *os.File) error {
if _, err := io.Copy(target, source); err != nil {
return fmt.Errorf("copy data: %w", err)
}
if err := target.Sync(); err != nil {
return fmt.Errorf("flush target: %w", err)
}
return nil
}
func extractedDirName(fileName string) string {
lowerName := strings.ToLower(fileName)
if strings.HasSuffix(lowerName, tarGzSuffix) {
return fileName[:len(fileName)-len(tarGzSuffix)]
}
return fileName
}
func snapshotKey(s fileSnapshot) string {
return fmt.Sprintf("%d:%d", s.size, s.modTime.UnixNano())
}
func isStableEnough(cfg config, stableCount int, snap fileSnapshot, now time.Time) bool {
if stableCount < cfg.stableScans {
return false
}
return now.Sub(snap.modTime) >= cfg.stableFor
}
func pathExists(path string) bool {
_, err := os.Stat(path)
return err == nil
}
func isDir(path string) bool {
info, err := os.Stat(path)
return err == nil && info.IsDir()
}