service.go 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761
  1. package storage
  2. /*
  3. ArozOS Cluster Storage (ACS)
  4. The storage service turns the metadata index into a usable file system:
  5. - volumes: folders on local drives that this node contributes
  6. - transfer: signed, chunked (4 MiB, SHA-256 verified) moves of whole files
  7. between nodes, sized for Cloudflare's request limits
  8. - placement: which volume receives a new file
  9. - reads and writes for the cluster:/ drive (mod/filesystem/abstractions/clusterfs)
  10. - reconcile: rescans of local volumes so real files stay authoritative
  11. Files are never split across nodes; a file is one ordinary file on one or
  12. more volumes.
  13. */
  14. import (
  15. "context"
  16. "crypto/sha256"
  17. "encoding/hex"
  18. "errors"
  19. "fmt"
  20. "io"
  21. "net/http"
  22. "os"
  23. "path/filepath"
  24. "sort"
  25. "strings"
  26. "sync"
  27. "time"
  28. uuid "github.com/satori/go.uuid"
  29. "imuslab.com/arozos/mod/cluster/membership"
  30. "imuslab.com/arozos/mod/cluster/metadata"
  31. "imuslab.com/arozos/mod/cluster/scheduling"
  32. "imuslab.com/arozos/mod/info/logger"
  33. )
  34. const (
  35. // ChunkSize is the transfer unit between nodes.
  36. ChunkSize = 4 << 20
  37. // MaxChunkSize bounds a single read request.
  38. MaxChunkSize = 8 << 20
  39. // SessionTTL is how long an unfinished upload session is kept.
  40. SessionTTL = 30 * time.Minute
  41. // placementHeadroom is kept free on a volume beyond the file size.
  42. placementHeadroom = 64 << 20
  43. )
  44. var (
  45. ErrNoVolume = errors.New("no cluster volume available; add one in System Settings > Cluster Settings > Storage")
  46. ErrNoHealthyCopy = errors.New("no healthy copy of this file is reachable right now")
  47. ErrNotLocalVolume = errors.New("volume is not on this node")
  48. ErrPathEscape = errors.New("path escapes the volume")
  49. ErrIsDirectory = errors.New("is a directory")
  50. ErrNotDirectory = errors.New("not a directory")
  51. ErrNotEmpty = errors.New("directory not empty")
  52. ErrChecksum = errors.New("checksum mismatch")
  53. )
  54. // Option configures the storage service.
  55. type Option struct {
  56. Membership *membership.Manager
  57. Metadata *metadata.Manager
  58. //Scheduler scores the nodes when choosing where a new file goes.
  59. Scheduler *scheduling.Manager
  60. TmpDir string
  61. // LocalRoots maps local (non network, non buffered) file system handler
  62. // UUIDs to their real root paths. Provided by the core.
  63. LocalRoots func() map[string]string
  64. // Thumbnailer renders the thumbnail of a file on one of this node's
  65. // volumes, given its OS path (thumbnail.go). Nil: this node renders none.
  66. Thumbnailer func(osPath string) ([]byte, error)
  67. RefreshInterval time.Duration
  68. ReconcileInterval time.Duration
  69. }
  70. // Service is the storage layer of this node.
  71. type Service struct {
  72. m *membership.Manager
  73. meta *metadata.Manager
  74. client *Client
  75. sched *scheduling.Manager
  76. opt Option
  77. sessMu sync.Mutex
  78. sessions map[string]*session
  79. //thumbSem bounds the thumbnails rendered at once
  80. thumbSem chan struct{}
  81. stop chan struct{}
  82. wg sync.WaitGroup
  83. once sync.Once
  84. //Hooks for the event bus (Phase 6); may be nil
  85. OnFileWritten func(rec metadata.FileRecord)
  86. OnFileRemoved func(rec metadata.FileRecord)
  87. OnFileRenamed func(oldPath string, rec metadata.FileRecord)
  88. //Replica hooks fired by PullCopy / VerifyLocalCopies
  89. OnReplicaVerified func(rec metadata.FileRecord, volumeID string)
  90. OnReplicaStale func(rec metadata.FileRecord, volumeID string)
  91. //OnDiskFull fires when a volume drops under the low water mark.
  92. OnDiskFull func(vol metadata.Volume)
  93. }
  94. // New creates the service and registers its node endpoints.
  95. func New(opt Option) (*Service, error) {
  96. if opt.Membership == nil || opt.Metadata == nil || opt.LocalRoots == nil {
  97. return nil, errors.New("membership, metadata and local roots are required")
  98. }
  99. if opt.TmpDir == "" {
  100. opt.TmpDir = os.TempDir()
  101. }
  102. if opt.RefreshInterval <= 0 {
  103. opt.RefreshInterval = 60 * time.Second
  104. }
  105. if opt.ReconcileInterval <= 0 {
  106. opt.ReconcileInterval = 30 * time.Minute
  107. }
  108. os.MkdirAll(filepath.Join(opt.TmpDir, "cluster"), 0755)
  109. if opt.Scheduler == nil {
  110. sc, err := scheduling.New(opt.Membership, opt.Metadata)
  111. if err != nil {
  112. return nil, err
  113. }
  114. opt.Scheduler = sc
  115. }
  116. s := &Service{
  117. m: opt.Membership,
  118. meta: opt.Metadata,
  119. sched: opt.Scheduler,
  120. client: &Client{Transport: opt.Membership.Transport()},
  121. opt: opt,
  122. sessions: map[string]*session{},
  123. thumbSem: make(chan struct{}, ThumbnailConcurrency),
  124. stop: make(chan struct{}),
  125. }
  126. s.registerACNHandlers()
  127. s.wg.Add(1)
  128. go s.maintenanceLoop()
  129. return s, nil
  130. }
  131. // Close stops background work.
  132. func (s *Service) Close() {
  133. s.once.Do(func() { close(s.stop) })
  134. s.wg.Wait()
  135. }
  136. // Ready reports whether the cluster:/ drive can serve requests.
  137. func (s *Service) Ready() bool {
  138. return s.m.InCluster()
  139. }
  140. func (s *Service) maintenanceLoop() {
  141. defer s.wg.Done()
  142. refresh := time.NewTicker(s.opt.RefreshInterval)
  143. reconcile := time.NewTicker(s.opt.ReconcileInterval)
  144. sweep := time.NewTicker(time.Minute)
  145. defer refresh.Stop()
  146. defer reconcile.Stop()
  147. defer sweep.Stop()
  148. //First reconcile shortly after boot picks up files already in contributed folders
  149. first := time.NewTimer(10 * time.Second)
  150. defer first.Stop()
  151. for {
  152. select {
  153. case <-s.stop:
  154. return
  155. case <-refresh.C:
  156. s.refreshVolumes()
  157. case <-first.C:
  158. s.ReconcileAll()
  159. case <-reconcile.C:
  160. s.ReconcileAll()
  161. case <-sweep.C:
  162. s.sweepSessions()
  163. }
  164. }
  165. }
  166. /*
  167. Placement
  168. */
  169. func (s *Service) nodeState(nodeID string) membership.NodeState {
  170. for _, n := range s.m.NodeViews() {
  171. if n.ID == nodeID {
  172. return n.State
  173. }
  174. }
  175. return membership.StateOffline
  176. }
  177. func usable(state membership.NodeState) bool {
  178. return state == membership.StateOnline || state == membership.StateDegraded
  179. }
  180. // pickVolume chooses the volume that should receive a new copy of size
  181. // bytes. Nodes are ranked by the cluster scorer (so a busy, distant or
  182. // nearly full node loses), preferNode and preferVolume win ties, and the
  183. // roomiest writable volume of the winning node is used.
  184. func (s *Service) pickVolume(size int64, preferNode string, preferVolume string, exclude []string) (*metadata.Volume, error) {
  185. need := size + size/10 + placementHeadroom
  186. excluded := map[string]bool{}
  187. for _, id := range exclude {
  188. excluded[id] = true
  189. }
  190. //Group the usable volumes by node
  191. byNode := map[string][]metadata.Volume{}
  192. for _, v := range s.meta.Volumes() {
  193. if v.Removed || v.ReadOnly || v.Evacuating || v.Free < need || excluded[v.ID] {
  194. continue
  195. }
  196. byNode[v.NodeID] = append(byNode[v.NodeID], v)
  197. }
  198. if len(byNode) == 0 {
  199. return nil, ErrNoVolume
  200. }
  201. candidates := []scheduling.Candidate{}
  202. for _, n := range s.m.NodeViews() {
  203. vols, ok := byNode[n.ID]
  204. if !ok {
  205. continue
  206. }
  207. best := 0.0
  208. for _, v := range vols {
  209. if v.Capacity > 0 {
  210. if f := float64(v.Free) / float64(v.Capacity); f > best {
  211. best = f
  212. }
  213. }
  214. }
  215. c := scheduling.Candidate{Node: n, FreeDisk: best, LatencyToData: -1, Diversity: -1}
  216. if n.ID == preferNode {
  217. c.Locality = 1 //writing where the bytes already are costs nothing
  218. }
  219. candidates = append(candidates, c)
  220. }
  221. ranked := s.sched.Rank(candidates, func(c scheduling.Candidate) (bool, string) {
  222. if !usable(c.Node.State) {
  223. return false, "node is " + strings.ToLower(string(c.Node.State))
  224. }
  225. return true, ""
  226. })
  227. nodeID, why := scheduling.Best(ranked)
  228. if nodeID == "" {
  229. if why != "" {
  230. //Keep the sentinel so callers can still test for it, but say
  231. //why the only candidate was turned down
  232. return nil, fmt.Errorf("%w (%s)", ErrNoVolume, why)
  233. }
  234. return nil, ErrNoVolume
  235. }
  236. vols := byNode[nodeID]
  237. sort.Slice(vols, func(i, j int) bool {
  238. if (vols[i].ID == preferVolume) != (vols[j].ID == preferVolume) {
  239. return vols[i].ID == preferVolume
  240. }
  241. if vols[i].Free != vols[j].Free {
  242. return vols[i].Free > vols[j].Free
  243. }
  244. return vols[i].ID < vols[j].ID
  245. })
  246. v := vols[0]
  247. return &v, nil
  248. }
  249. // PlaceRequest asks the leader where a new file should go.
  250. type PlaceRequest struct {
  251. Size int64 `json:"size"`
  252. Path string `json:"path"`
  253. PreferNode string `json:"preferNode"`
  254. PreferVolume string `json:"preferVolume"`
  255. ExcludeVolume string `json:"excludeVolume,omitempty"`
  256. }
  257. type PlaceResponse struct {
  258. VolumeID string `json:"volumeId"`
  259. }
  260. // place decides the volume for a write, consulting the leader when this
  261. // node is not the leader and the leader is reachable.
  262. func (s *Service) place(size int64, logical string, preferVolume string) (*metadata.Volume, error) {
  263. me := s.m.NodeID()
  264. if leader := s.meta.Leader(); leader != "" && leader != me {
  265. ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
  266. defer cancel()
  267. var resp PlaceResponse
  268. err := s.m.Transport().DoJSON(ctx, leader, http.MethodPost, pathPlace, PlaceRequest{Size: size, Path: logical, PreferNode: me, PreferVolume: preferVolume}, &resp)
  269. if err == nil && resp.VolumeID != "" {
  270. if v, ok := s.meta.Volume(resp.VolumeID); ok && !v.Removed {
  271. return v, nil
  272. }
  273. }
  274. }
  275. return s.pickVolume(size, me, preferVolume, nil)
  276. }
  277. /*
  278. Backend for the cluster:/ drive
  279. */
  280. // Info is the file-system view of a record.
  281. type Info struct {
  282. Path string
  283. IsDir bool
  284. Size int64
  285. ModTime int64
  286. }
  287. func infoOf(rec *metadata.FileRecord) Info {
  288. return Info{Path: rec.Path, IsDir: rec.IsDir, Size: rec.Size, ModTime: rec.ModTime}
  289. }
  290. // Stat returns the entry at a logical path.
  291. func (s *Service) Stat(logical string) (Info, error) {
  292. logical = metadata.NormalizePath(logical)
  293. if logical == "/" {
  294. return Info{Path: "/", IsDir: true, ModTime: time.Now().Unix()}, nil
  295. }
  296. rec, err := s.meta.Stat(logical)
  297. if err != nil {
  298. return Info{}, os.ErrNotExist
  299. }
  300. return infoOf(rec), nil
  301. }
  302. // List returns the children of a directory.
  303. func (s *Service) List(logical string) ([]Info, error) {
  304. logical = metadata.NormalizePath(logical)
  305. if logical != "/" {
  306. rec, err := s.meta.Stat(logical)
  307. if err != nil {
  308. return nil, os.ErrNotExist
  309. }
  310. if !rec.IsDir {
  311. return nil, ErrNotDirectory
  312. }
  313. }
  314. recs, err := s.meta.ListDir(logical)
  315. if err != nil {
  316. return nil, err
  317. }
  318. out := make([]Info, 0, len(recs))
  319. for i := range recs {
  320. out = append(out, infoOf(&recs[i]))
  321. }
  322. return out, nil
  323. }
  324. // ensureDirs creates directory records for every missing parent of logical.
  325. func (s *Service) ensureDirs(logical string, owner string) error {
  326. logical = metadata.NormalizePath(logical)
  327. if logical == "/" {
  328. return nil
  329. }
  330. parent := metadata.ParentPath(logical)
  331. if err := s.ensureDirs(parent, owner); err != nil {
  332. return err
  333. }
  334. if rec, err := s.meta.Stat(logical); err == nil {
  335. if !rec.IsDir {
  336. return ErrNotDirectory
  337. }
  338. return nil
  339. }
  340. return s.meta.Submit(metadata.KindFile, &metadata.FileRecord{
  341. ID: uuid.NewV4().String(),
  342. Path: logical,
  343. IsDir: true,
  344. ModTime: time.Now().Unix(),
  345. Owner: owner,
  346. })
  347. }
  348. // Mkdir creates a directory (and its parents) in the namespace.
  349. func (s *Service) Mkdir(logical string, owner string) error {
  350. if !s.Ready() {
  351. return metadata.ErrNotInCluster
  352. }
  353. logical = metadata.NormalizePath(logical)
  354. if rec, err := s.meta.Stat(logical); err == nil {
  355. if rec.IsDir {
  356. return nil
  357. }
  358. return metadata.ErrExists
  359. }
  360. return s.ensureDirs(logical, owner)
  361. }
  362. // Write stores a new version of a file: spool + hash locally, place, copy,
  363. // then publish the record.
  364. func (s *Service) Write(logical string, r io.Reader, owner string) error {
  365. if !s.Ready() {
  366. return metadata.ErrNotInCluster
  367. }
  368. logical = metadata.NormalizePath(logical)
  369. if logical == "/" {
  370. return ErrIsDirectory
  371. }
  372. if rec, err := s.meta.Stat(logical); err == nil && rec.IsDir {
  373. return ErrIsDirectory
  374. }
  375. if err := s.ensureDirs(metadata.ParentPath(logical), owner); err != nil {
  376. return err
  377. }
  378. //1. spool and hash
  379. spool, err := os.CreateTemp(filepath.Join(s.opt.TmpDir, "cluster"), "spool-*")
  380. if err != nil {
  381. return err
  382. }
  383. spoolPath := spool.Name()
  384. defer os.Remove(spoolPath)
  385. h := sha256.New()
  386. size, err := io.Copy(io.MultiWriter(spool, h), r)
  387. spool.Close()
  388. if err != nil {
  389. return err
  390. }
  391. checksum := hex.EncodeToString(h.Sum(nil))
  392. //2. placement (overwrite prefers the existing primary)
  393. preferVolume := ""
  394. existing, _ := s.meta.Stat(logical)
  395. if existing != nil {
  396. preferVolume = existing.Primary
  397. }
  398. vol, err := s.place(size, logical, preferVolume)
  399. if err != nil {
  400. return err
  401. }
  402. //3. copy bytes first: the namespace only ever shows committed files
  403. fileID := uuid.NewV4().String()
  404. if existing != nil {
  405. fileID = existing.ID
  406. }
  407. if vol.NodeID == s.m.NodeID() {
  408. err = s.writeLocalFile(vol, logical, spoolPath, size, checksum)
  409. } else {
  410. err = s.uploadRemote(vol, logical, fileID, spoolPath, size, checksum)
  411. }
  412. if err != nil {
  413. return err
  414. }
  415. //4. publish the record with the new copy as its only location
  416. rec := &metadata.FileRecord{ID: fileID, Path: logical, Owner: owner}
  417. if existing != nil {
  418. rec = existing.Clone()
  419. rec.Locations = nil
  420. }
  421. rec.IsDir = false
  422. rec.Size = size
  423. rec.Checksum = checksum
  424. rec.ModTime = time.Now().Unix()
  425. rec.Primary = vol.ID
  426. rec.SetLocation(metadata.Location{VolumeID: vol.ID, NodeID: vol.NodeID, State: metadata.LocCommitted, Checksum: checksum})
  427. if err := s.meta.Submit(metadata.KindFile, rec); err != nil {
  428. return err
  429. }
  430. if existing != nil {
  431. //Physically remove superseded copies on other volumes (best effort)
  432. for _, old := range existing.Locations {
  433. if old.VolumeID != vol.ID {
  434. s.deletePhysical(old, logical)
  435. }
  436. }
  437. }
  438. if s.OnFileWritten != nil {
  439. s.OnFileWritten(*rec)
  440. }
  441. return nil
  442. }
  443. func (s *Service) writeLocalFile(vol *metadata.Volume, logical string, spoolPath string, size int64, checksum string) error {
  444. real, err := s.realPath(vol, logical)
  445. if err != nil {
  446. return err
  447. }
  448. if err := os.MkdirAll(filepath.Dir(real), 0755); err != nil {
  449. return err
  450. }
  451. part := real + ".part-" + uuid.NewV4().String()
  452. if err := copyFile(spoolPath, part); err != nil {
  453. os.Remove(part)
  454. return err
  455. }
  456. if err := os.Rename(part, real); err != nil {
  457. os.Remove(part)
  458. return err
  459. }
  460. return nil
  461. }
  462. func (s *Service) uploadRemote(vol *metadata.Volume, logical string, fileID string, spoolPath string, size int64, checksum string) error {
  463. f, err := os.Open(spoolPath)
  464. if err != nil {
  465. return err
  466. }
  467. defer f.Close()
  468. ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour)
  469. defer cancel()
  470. return s.client.Upload(ctx, vol.NodeID, vol.ID, logical, fileID, f, size, checksum, nil)
  471. }
  472. // OpenRead returns a stream of the file from the best available copy.
  473. func (s *Service) OpenRead(logical string) (io.ReadCloser, error) {
  474. if !s.Ready() {
  475. return nil, metadata.ErrNotInCluster
  476. }
  477. rec, err := s.meta.Stat(logical)
  478. if err != nil {
  479. return nil, os.ErrNotExist
  480. }
  481. if rec.IsDir {
  482. return nil, ErrIsDirectory
  483. }
  484. locs := s.orderedLocations(rec)
  485. if len(locs) == 0 {
  486. return nil, ErrNoHealthyCopy
  487. }
  488. var lastErr error = ErrNoHealthyCopy
  489. for _, loc := range locs {
  490. vol, ok := s.meta.Volume(loc.VolumeID)
  491. if !ok || vol.Removed {
  492. continue
  493. }
  494. if vol.NodeID == s.m.NodeID() {
  495. real, err := s.realPath(vol, rec.Path)
  496. if err != nil {
  497. lastErr = err
  498. continue
  499. }
  500. f, err := os.Open(real)
  501. if err != nil {
  502. lastErr = err
  503. continue
  504. }
  505. return f, nil
  506. }
  507. //Remote copy: stream through a pipe while verifying the checksum
  508. pr, pw := io.Pipe()
  509. go func(v metadata.Volume) {
  510. ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour)
  511. defer cancel()
  512. err := s.client.Download(ctx, v.NodeID, v.ID, rec.Path, pw, rec.Checksum, nil)
  513. pw.CloseWithError(err)
  514. }(*vol)
  515. return pr, nil
  516. }
  517. return nil, lastErr
  518. }
  519. // orderedLocations lists readable copies, local first, then online nodes in
  520. // membership order.
  521. func (s *Service) orderedLocations(rec *metadata.FileRecord) []metadata.Location {
  522. me := s.m.NodeID()
  523. out := []metadata.Location{}
  524. for _, l := range rec.HealthyLocations() {
  525. if l.NodeID == me {
  526. out = append(out, l)
  527. }
  528. }
  529. for _, l := range rec.HealthyLocations() {
  530. if l.NodeID != me && usable(s.nodeState(l.NodeID)) {
  531. out = append(out, l)
  532. }
  533. }
  534. return out
  535. }
  536. // Remove deletes a file or directory from the namespace and its copies.
  537. func (s *Service) Remove(logical string, recursive bool) error {
  538. if !s.Ready() {
  539. return metadata.ErrNotInCluster
  540. }
  541. logical = metadata.NormalizePath(logical)
  542. if logical == "/" {
  543. return errors.New("cannot remove the namespace root")
  544. }
  545. rec, err := s.meta.Stat(logical)
  546. if err != nil {
  547. return os.ErrNotExist
  548. }
  549. targets := []metadata.FileRecord{*rec}
  550. if rec.IsDir {
  551. children := s.meta.ListSubtree(logical)
  552. if len(children) > 0 && !recursive {
  553. return ErrNotEmpty
  554. }
  555. targets = append(targets, children...)
  556. }
  557. for i := range targets {
  558. t := targets[i].Clone()
  559. t.Removed = true
  560. if err := s.meta.Submit(metadata.KindFile, t); err != nil {
  561. return err
  562. }
  563. if !t.IsDir {
  564. for _, loc := range t.Locations {
  565. s.deletePhysical(loc, t.Path)
  566. }
  567. } else {
  568. s.removePhysicalDir(t.Path)
  569. }
  570. if s.OnFileRemoved != nil {
  571. s.OnFileRemoved(*t)
  572. }
  573. }
  574. return nil
  575. }
  576. // Rename moves a file or directory (and everything beneath it).
  577. func (s *Service) Rename(oldPath string, newPath string) error {
  578. if !s.Ready() {
  579. return metadata.ErrNotInCluster
  580. }
  581. oldPath = metadata.NormalizePath(oldPath)
  582. newPath = metadata.NormalizePath(newPath)
  583. if oldPath == "/" || newPath == "/" || oldPath == newPath {
  584. return errors.New("invalid rename")
  585. }
  586. rec, err := s.meta.Stat(oldPath)
  587. if err != nil {
  588. return os.ErrNotExist
  589. }
  590. if _, err := s.meta.Stat(newPath); err == nil {
  591. return metadata.ErrExists
  592. }
  593. if err := s.ensureDirs(metadata.ParentPath(newPath), rec.Owner); err != nil {
  594. return err
  595. }
  596. targets := []metadata.FileRecord{*rec}
  597. if rec.IsDir {
  598. targets = append(targets, s.meta.ListSubtree(oldPath)...)
  599. }
  600. for i := range targets {
  601. t := targets[i].Clone()
  602. from := t.Path
  603. t.Path = newPath + t.Path[len(oldPath):]
  604. if !t.IsDir {
  605. for _, loc := range t.Locations {
  606. if err := s.renamePhysical(loc, from, t.Path); err != nil {
  607. loc.State = metadata.LocStale
  608. t.SetLocation(loc)
  609. }
  610. }
  611. }
  612. if err := s.meta.Submit(metadata.KindFile, t); err != nil {
  613. return err
  614. }
  615. if s.OnFileRenamed != nil {
  616. s.OnFileRenamed(from, *t)
  617. }
  618. }
  619. s.removePhysicalDir(oldPath)
  620. return nil
  621. }
  622. /*
  623. Physical helpers (local or through the transfer client)
  624. */
  625. func (s *Service) deletePhysical(loc metadata.Location, logical string) {
  626. vol, ok := s.meta.Volume(loc.VolumeID)
  627. if !ok {
  628. return
  629. }
  630. if vol.NodeID == s.m.NodeID() {
  631. if real, err := s.realPath(vol, logical); err == nil {
  632. os.Remove(real)
  633. }
  634. return
  635. }
  636. ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
  637. defer cancel()
  638. if err := s.client.Delete(ctx, vol.NodeID, vol.ID, logical, false); err != nil {
  639. logger.PrintAndLog("Cluster", "Unable to delete copy of "+logical+" on "+s.m.NodeName(vol.NodeID)+": "+err.Error(), nil)
  640. }
  641. }
  642. // removePhysicalDir removes now-empty directories on every volume (best effort).
  643. func (s *Service) removePhysicalDir(logical string) {
  644. for _, vol := range s.meta.Volumes() {
  645. if vol.Removed {
  646. continue
  647. }
  648. if vol.NodeID == s.m.NodeID() {
  649. if real, err := s.realPath(&vol, logical); err == nil {
  650. os.Remove(real) //fails when not empty, which is fine
  651. }
  652. continue
  653. }
  654. if usable(s.nodeState(vol.NodeID)) {
  655. ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
  656. s.client.Delete(ctx, vol.NodeID, vol.ID, logical, false)
  657. cancel()
  658. }
  659. }
  660. }
  661. func (s *Service) renamePhysical(loc metadata.Location, from string, to string) error {
  662. vol, ok := s.meta.Volume(loc.VolumeID)
  663. if !ok {
  664. return ErrNotLocalVolume
  665. }
  666. if vol.NodeID == s.m.NodeID() {
  667. src, err := s.realPath(vol, from)
  668. if err != nil {
  669. return err
  670. }
  671. dst, err := s.realPath(vol, to)
  672. if err != nil {
  673. return err
  674. }
  675. if err := os.MkdirAll(filepath.Dir(dst), 0755); err != nil {
  676. return err
  677. }
  678. return os.Rename(src, dst)
  679. }
  680. ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
  681. defer cancel()
  682. return s.client.Rename(ctx, vol.NodeID, vol.ID, from, to)
  683. }
  684. func copyFile(src string, dst string) error {
  685. in, err := os.Open(src)
  686. if err != nil {
  687. return err
  688. }
  689. defer in.Close()
  690. out, err := os.Create(dst)
  691. if err != nil {
  692. return err
  693. }
  694. if _, err := io.Copy(out, in); err != nil {
  695. out.Close()
  696. return err
  697. }
  698. return out.Close()
  699. }
  700. func fileChecksum(path string) (string, int64, error) {
  701. f, err := os.Open(path)
  702. if err != nil {
  703. return "", 0, err
  704. }
  705. defer f.Close()
  706. h := sha256.New()
  707. n, err := io.Copy(h, f)
  708. if err != nil {
  709. return "", 0, err
  710. }
  711. return hex.EncodeToString(h.Sum(nil)), n, nil
  712. }