@@ -84,6 +84,16 @@ type fsRestore struct {
8484 apfl * pgalloc.AsyncPagesFileLoad
8585 mfs map [checkpoint.ResourceID ]* fscheckpoint.MemoryFile
8686 tmpfs map [checkpoint.ResourceID ]* fscheckpoint.Tmpfs
87+
88+ waitMu sync.Mutex
89+ waitMap map [string ]* fsRestoreContainer // key is container ID
90+ }
91+
92+ type fsRestoreContainer struct {
93+ // These fields are protected by fsRestore.waitMu.
94+ err error
95+ asyncLoads int // number of MemoryFiles currently in async page loading
96+ cond sync.Cond
8797}
8898
8999// fsRestoreOpts holds options to startFSRestore.
@@ -111,8 +121,9 @@ func makeFSRestoreOptsForLocalCheckpoint(args *Args) (fsRestoreOpts, error) {
111121// startFSRestore takes ownership of resources in opts.
112122func startFSRestore (opts * fsRestoreOpts ) (* fsRestore , error ) {
113123 fsr := & fsRestore {
114- mfs : make (map [checkpoint.ResourceID ]* fscheckpoint.MemoryFile ),
115- tmpfs : make (map [checkpoint.ResourceID ]* fscheckpoint.Tmpfs ),
124+ mfs : make (map [checkpoint.ResourceID ]* fscheckpoint.MemoryFile ),
125+ tmpfs : make (map [checkpoint.ResourceID ]* fscheckpoint.Tmpfs ),
126+ waitMap : make (map [string ]* fsRestoreContainer ),
116127 }
117128
118129 // TODO: NOLINT - Currently we read the whole pages metadata file into a
@@ -225,32 +236,71 @@ func startFSRestore(opts *fsRestoreOpts) (*fsRestore, error) {
225236 return fsr , nil
226237}
227238
228- func (fsr * fsRestore ) memoryFileLoadArgs (id checkpoint.ResourceID ) (io.Reader , uint64 , error ) {
239+ // +checklocks:fsr.waitMu
240+ func (fsr * fsRestore ) ensureContainer (cid string ) * fsRestoreContainer {
241+ c := fsr .waitMap [cid ]
242+ if c == nil {
243+ c = & fsRestoreContainer {}
244+ c .cond .L = & fsr .waitMu
245+ fsr .waitMap [cid ] = c
246+ }
247+ return c
248+ }
249+
250+ func (c * fsRestoreContainer ) setError (err error ) error {
251+ if c .err == nil && err != nil {
252+ c .err = err
253+ c .cond .Broadcast ()
254+ }
255+ return err
256+ }
257+
258+ func (fsr * fsRestore ) memoryFileLoadArgs (id checkpoint.ResourceID , cid string ) (io.Reader , uint64 , func (error ), error ) {
229259 if fsr == nil {
230- return nil , 0 , nil
260+ return nil , 0 , func ( error ) {}, nil
231261 }
262+
232263 fsr .wg .Wait ()
233264 if fsr .manifestErr != nil {
234- return nil , 0 , fsr .manifestErr
265+ return nil , 0 , nil , fsr .manifestErr
235266 }
236267 mmf := fsr .mfs [id ]
237268 if mmf == nil {
238- return nil , 0 , nil
269+ return nil , 0 , func ( error ) {}, nil
239270 }
240271 pagesMetadata , err := fsr .getPagesMetadata ()
241- if mmf .PagesMetadataEnd <= uint64 (len (pagesMetadata )) {
242- return bytes .NewReader (pagesMetadata [mmf .PagesMetadataStart :mmf .PagesMetadataEnd ]), mmf .PagesStart , nil
243- }
244- if err != nil {
245- return nil , 0 , fmt .Errorf ("failed to read pages metadata: %w" , err )
272+
273+ fsr .waitMu .Lock ()
274+ defer fsr .waitMu .Unlock ()
275+ c := fsr .ensureContainer (cid )
276+ if mmf .PagesMetadataEnd > uint64 (len (pagesMetadata )) {
277+ if err != nil {
278+ return nil , 0 , nil , c .setError (fmt .Errorf ("failed to read pages metadata: %w" , err ))
279+ }
280+ return nil , 0 , nil , c .setError (fmt .Errorf ("MemoryFile %q has pages metadata range [%d, %d) beyond pages metadata file size %d" , mmf .ResourceID , mmf .PagesMetadataStart , mmf .PagesMetadataEnd , len (pagesMetadata )))
246281 }
247- return nil , 0 , fmt .Errorf ("MemoryFile %q has pages metadata range [%d, %d) beyond pages metadata file size %d" , mmf .ResourceID , mmf .PagesMetadataStart , mmf .PagesMetadataEnd , len (pagesMetadata ))
282+ c .asyncLoads ++
283+ return bytes .NewReader (pagesMetadata [mmf .PagesMetadataStart :mmf .PagesMetadataEnd ]), mmf .PagesStart , func (err error ) {
284+ fsr .waitMu .Lock ()
285+ defer fsr .waitMu .Unlock ()
286+ c .asyncLoads --
287+ switch {
288+ case err != nil :
289+ if c .err == nil {
290+ c .err = err
291+ }
292+ fallthrough
293+ case c .asyncLoads == 0 :
294+ c .cond .Broadcast ()
295+ }
296+ }, nil
248297}
249298
250- func (fsr * fsRestore ) tmpfsSourceTar (id checkpoint.ResourceID ) (io.ReadCloser , error ) {
299+ func (fsr * fsRestore ) tmpfsSourceTar (id checkpoint.ResourceID , cid string ) (io.ReadCloser , error ) {
251300 if fsr == nil {
252301 return nil , nil
253302 }
303+
254304 fsr .wg .Wait ()
255305 if fsr .manifestErr != nil {
256306 return nil , fsr .manifestErr
@@ -260,11 +310,43 @@ func (fsr *fsRestore) tmpfsSourceTar(id checkpoint.ResourceID) (io.ReadCloser, e
260310 return nil , nil
261311 }
262312 multiTar , err := fsr .getMultiTar ()
313+
314+ fsr .waitMu .Lock ()
315+ defer fsr .waitMu .Unlock ()
263316 if mt .TarEnd <= uint64 (len (multiTar )) {
264317 return io .NopCloser (bytes .NewReader (multiTar [mt .TarStart :mt .TarEnd ])), nil
265318 }
319+ c := fsr .ensureContainer (cid )
266320 if err != nil {
267- return nil , fmt .Errorf ("failed to read tar archive: %w" , err )
321+ return nil , c .setError (fmt .Errorf ("failed to read tar archive: %w" , err ))
322+ }
323+ return nil , c .setError (fmt .Errorf ("tmpfs %q has tar range [%d, %d) beyond multi-tar file size %d" , mt .ResourceID , mt .TarStart , mt .TarEnd , len (multiTar )))
324+ }
325+
326+ // wait blocks until either all filesystems have been restored for the
327+ // container with the given ID, or an error occurs while restoring filesystems
328+ // for that container.
329+ func (fsr * fsRestore ) wait (cid string ) error {
330+ if fsr == nil {
331+ return fmt .Errorf ("filesystem restore is not enabled" )
332+ }
333+ fsr .wg .Wait ()
334+ if fsr .manifestErr != nil {
335+ return fsr .manifestErr
336+ }
337+ fsr .waitMu .Lock ()
338+ defer fsr .waitMu .Unlock ()
339+ c := fsr .waitMap [cid ]
340+ if c == nil {
341+ return fmt .Errorf ("no filesystems restored for container %s" , cid )
342+ }
343+ for {
344+ if c .err != nil {
345+ return c .err
346+ }
347+ if c .asyncLoads == 0 {
348+ return nil
349+ }
350+ c .cond .Wait ()
268351 }
269- return nil , fmt .Errorf ("tmpfs %q has tar range [%d, %d) beyond multi-tar file size %d" , mt .ResourceID , mt .TarStart , mt .TarEnd , len (multiTar ))
270352}
0 commit comments