55 "crypto/sha256"
66 "errors"
77 "fmt"
8- "runtime"
9- "sync"
108 "sync/atomic"
119
1210 "github.com/tetratelabs/wazero"
@@ -199,54 +197,6 @@ func (r *Runtime) Acquire(ctx context.Context) (*Instance, error) {
199197 }
200198}
201199
202- func (r * Runtime ) AcquireN (ctx context.Context , n int ) ([]* Instance , error ) {
203- if n <= 0 {
204- return nil , nil
205- }
206- out := make ([]* Instance , 0 , n )
207- drained := 0
208- for len (out ) < n {
209- select {
210- case inst := <- r .warm :
211- out = append (out , inst )
212- drained ++
213- default :
214- goto done
215- }
216- }
217- done:
218- if drained > 0 {
219- refillCount := drained
220- go func () {
221- refills , err := r .RestoreN (context .Background (), r .golden , refillCount )
222- if err != nil {
223- for _ , inst := range refills {
224- inst .Release ()
225- }
226- return
227- }
228- for _ , inst := range refills {
229- select {
230- case r .warm <- inst :
231- default :
232- inst .Release ()
233- }
234- }
235- }()
236- }
237- if remaining := n - len (out ); remaining > 0 {
238- rest , err := r .RestoreN (ctx , r .golden , remaining )
239- if err != nil {
240- for _ , inst := range out {
241- inst .Release ()
242- }
243- return nil , err
244- }
245- out = append (out , rest ... )
246- }
247- return out , nil
248- }
249-
250200func (r * Runtime ) Restore (ctx context.Context , s Snapshot ) (* Instance , error ) {
251201 adapterID , hash , memory , err := decodeSnapshot (s )
252202 if err != nil {
@@ -258,73 +208,6 @@ func (r *Runtime) Restore(ctx context.Context, s Snapshot) (*Instance, error) {
258208 return r .restoreDecoded (ctx , memory )
259209}
260210
261- func (r * Runtime ) RestoreN (ctx context.Context , s Snapshot , n int ) ([]* Instance , error ) {
262- if n <= 0 {
263- return nil , nil
264- }
265- adapterID , hash , memory , err := decodeSnapshot (s )
266- if err != nil {
267- return nil , err
268- }
269- if err := r .checkSnapshotIdentity (adapterID , hash ); err != nil {
270- return nil , err
271- }
272-
273- if n == 1 {
274- inst , err := r .restoreDecoded (ctx , memory )
275- if err != nil {
276- return nil , err
277- }
278- return []* Instance {inst }, nil
279- }
280-
281- parallelism := r .forkParallelism
282- if parallelism <= 0 {
283- parallelism = runtime .GOMAXPROCS (0 )
284- }
285- if parallelism > n {
286- parallelism = n
287- }
288-
289- out := make ([]* Instance , n )
290- sem := make (chan struct {}, parallelism )
291- var wg sync.WaitGroup
292- var firstErr atomic.Value
293- cctx , cancel := context .WithCancel (ctx )
294- defer cancel ()
295-
296- for i := 0 ; i < n ; i ++ {
297- wg .Add (1 )
298- sem <- struct {}{}
299- go func (idx int ) {
300- defer wg .Done ()
301- defer func () { <- sem }()
302- if cctx .Err () != nil {
303- return
304- }
305- inst , err := r .restoreDecoded (cctx , memory )
306- if err != nil {
307- if firstErr .CompareAndSwap (nil , err ) {
308- cancel ()
309- }
310- return
311- }
312- out [idx ] = inst
313- }(i )
314- }
315- wg .Wait ()
316-
317- if v := firstErr .Load (); v != nil {
318- for _ , inst := range out {
319- if inst != nil {
320- inst .Release ()
321- }
322- }
323- return nil , v .(error )
324- }
325- return out , nil
326- }
327-
328211func (r * Runtime ) checkSnapshotIdentity (adapterID string , hash [32 ]byte ) error {
329212 if adapterID != r .adapter .ID () {
330213 return fmt .Errorf ("sango: snapshot adapter %q does not match runtime adapter %q" ,
0 commit comments