
There were races with the way process states. This displayed in ways, especially around pausing the container for atomic operations. Users would get errors like, cannnot delete container in paused state and such. This can be eaisly reproduced with `docker` and the following command: ```bash > (for i in `seq 1 25`; do id=$(docker create alpine usleep 50000);docker start $id;docker commit $id;docker wait $id;docker rm $id; done) ``` This two issues that this fixes are: * locks must be held by the owning process, not the state operations. * If a container ends up being paused but before the operation completes, the process exists, make sure we resume the container before setting the the process as exited. Signed-off-by: Michael Crosby <crosbymichael@gmail.com>
420 lines
11 KiB
Go
420 lines
11 KiB
Go
// +build !windows
|
|
|
|
/*
|
|
Copyright The containerd Authors.
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
you may not use this file except in compliance with the License.
|
|
You may obtain a copy of the License at
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
See the License for the specific language governing permissions and
|
|
limitations under the License.
|
|
*/
|
|
|
|
package proc
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"syscall"
|
|
|
|
"github.com/containerd/console"
|
|
"github.com/containerd/containerd/errdefs"
|
|
"github.com/containerd/containerd/runtime/proc"
|
|
"github.com/containerd/fifo"
|
|
runc "github.com/containerd/go-runc"
|
|
google_protobuf "github.com/gogo/protobuf/types"
|
|
"github.com/pkg/errors"
|
|
"github.com/sirupsen/logrus"
|
|
)
|
|
|
|
type initState interface {
|
|
Resize(console.WinSize) error
|
|
Start(context.Context) error
|
|
Delete(context.Context) error
|
|
Pause(context.Context) error
|
|
Resume(context.Context) error
|
|
Update(context.Context, *google_protobuf.Any) error
|
|
Checkpoint(context.Context, *CheckpointConfig) error
|
|
Exec(context.Context, string, *ExecConfig) (proc.Process, error)
|
|
Kill(context.Context, uint32, bool) error
|
|
SetExited(int)
|
|
}
|
|
|
|
type createdState struct {
|
|
p *Init
|
|
}
|
|
|
|
func (s *createdState) transition(name string) error {
|
|
switch name {
|
|
case "running":
|
|
s.p.initState = &runningState{p: s.p}
|
|
case "stopped":
|
|
s.p.initState = &stoppedState{p: s.p}
|
|
case "deleted":
|
|
s.p.initState = &deletedState{}
|
|
default:
|
|
return errors.Errorf("invalid state transition %q to %q", stateName(s), name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *createdState) Pause(ctx context.Context) error {
|
|
return errors.Errorf("cannot pause task in created state")
|
|
}
|
|
|
|
func (s *createdState) Resume(ctx context.Context) error {
|
|
return errors.Errorf("cannot resume task in created state")
|
|
}
|
|
|
|
func (s *createdState) Update(ctx context.Context, r *google_protobuf.Any) error {
|
|
return s.p.update(ctx, r)
|
|
}
|
|
|
|
func (s *createdState) Checkpoint(ctx context.Context, r *CheckpointConfig) error {
|
|
return errors.Errorf("cannot checkpoint a task in created state")
|
|
}
|
|
|
|
func (s *createdState) Resize(ws console.WinSize) error {
|
|
return s.p.resize(ws)
|
|
}
|
|
|
|
func (s *createdState) Start(ctx context.Context) error {
|
|
if err := s.p.start(ctx); err != nil {
|
|
return err
|
|
}
|
|
return s.transition("running")
|
|
}
|
|
|
|
func (s *createdState) Delete(ctx context.Context) error {
|
|
if err := s.p.delete(ctx); err != nil {
|
|
return err
|
|
}
|
|
return s.transition("deleted")
|
|
}
|
|
|
|
func (s *createdState) Kill(ctx context.Context, sig uint32, all bool) error {
|
|
return s.p.kill(ctx, sig, all)
|
|
}
|
|
|
|
func (s *createdState) SetExited(status int) {
|
|
s.p.setExited(status)
|
|
|
|
if err := s.transition("stopped"); err != nil {
|
|
panic(err)
|
|
}
|
|
}
|
|
|
|
func (s *createdState) Exec(ctx context.Context, path string, r *ExecConfig) (proc.Process, error) {
|
|
return s.p.exec(ctx, path, r)
|
|
}
|
|
|
|
type createdCheckpointState struct {
|
|
p *Init
|
|
opts *runc.RestoreOpts
|
|
}
|
|
|
|
func (s *createdCheckpointState) transition(name string) error {
|
|
switch name {
|
|
case "running":
|
|
s.p.initState = &runningState{p: s.p}
|
|
case "stopped":
|
|
s.p.initState = &stoppedState{p: s.p}
|
|
case "deleted":
|
|
s.p.initState = &deletedState{}
|
|
default:
|
|
return errors.Errorf("invalid state transition %q to %q", stateName(s), name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *createdCheckpointState) Pause(ctx context.Context) error {
|
|
return errors.Errorf("cannot pause task in created state")
|
|
}
|
|
|
|
func (s *createdCheckpointState) Resume(ctx context.Context) error {
|
|
return errors.Errorf("cannot resume task in created state")
|
|
}
|
|
|
|
func (s *createdCheckpointState) Update(ctx context.Context, r *google_protobuf.Any) error {
|
|
return s.p.update(ctx, r)
|
|
}
|
|
|
|
func (s *createdCheckpointState) Checkpoint(ctx context.Context, r *CheckpointConfig) error {
|
|
return errors.Errorf("cannot checkpoint a task in created state")
|
|
}
|
|
|
|
func (s *createdCheckpointState) Resize(ws console.WinSize) error {
|
|
return s.p.resize(ws)
|
|
}
|
|
|
|
func (s *createdCheckpointState) Start(ctx context.Context) error {
|
|
p := s.p
|
|
sio := p.stdio
|
|
|
|
var (
|
|
err error
|
|
socket *runc.Socket
|
|
)
|
|
if sio.Terminal {
|
|
if socket, err = runc.NewTempConsoleSocket(); err != nil {
|
|
return errors.Wrap(err, "failed to create OCI runtime console socket")
|
|
}
|
|
defer socket.Close()
|
|
s.opts.ConsoleSocket = socket
|
|
}
|
|
|
|
if _, err := s.p.runtime.Restore(ctx, p.id, p.Bundle, s.opts); err != nil {
|
|
return p.runtimeError(err, "OCI runtime restore failed")
|
|
}
|
|
if sio.Stdin != "" {
|
|
sc, err := fifo.OpenFifo(ctx, sio.Stdin, syscall.O_WRONLY|syscall.O_NONBLOCK, 0)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to open stdin fifo %s", sio.Stdin)
|
|
}
|
|
p.stdin = sc
|
|
p.closers = append(p.closers, sc)
|
|
}
|
|
var copyWaitGroup sync.WaitGroup
|
|
if socket != nil {
|
|
console, err := socket.ReceiveMaster()
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to retrieve console master")
|
|
}
|
|
console, err = p.Platform.CopyConsole(ctx, console, sio.Stdin, sio.Stdout, sio.Stderr, &p.wg, ©WaitGroup)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to start console copy")
|
|
}
|
|
p.console = console
|
|
} else if !sio.IsNull() {
|
|
if err := copyPipes(ctx, p.io, sio.Stdin, sio.Stdout, sio.Stderr, &p.wg, ©WaitGroup); err != nil {
|
|
return errors.Wrap(err, "failed to start io pipe copy")
|
|
}
|
|
}
|
|
|
|
copyWaitGroup.Wait()
|
|
pid, err := runc.ReadPidFile(s.opts.PidFile)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to retrieve OCI runtime container pid")
|
|
}
|
|
p.pid = pid
|
|
return s.transition("running")
|
|
}
|
|
|
|
func (s *createdCheckpointState) Delete(ctx context.Context) error {
|
|
if err := s.p.delete(ctx); err != nil {
|
|
return err
|
|
}
|
|
return s.transition("deleted")
|
|
}
|
|
|
|
func (s *createdCheckpointState) Kill(ctx context.Context, sig uint32, all bool) error {
|
|
return s.p.kill(ctx, sig, all)
|
|
}
|
|
|
|
func (s *createdCheckpointState) SetExited(status int) {
|
|
s.p.setExited(status)
|
|
|
|
if err := s.transition("stopped"); err != nil {
|
|
panic(err)
|
|
}
|
|
}
|
|
|
|
func (s *createdCheckpointState) Exec(ctx context.Context, path string, r *ExecConfig) (proc.Process, error) {
|
|
return nil, errors.Errorf("cannot exec in a created state")
|
|
}
|
|
|
|
type runningState struct {
|
|
p *Init
|
|
}
|
|
|
|
func (s *runningState) transition(name string) error {
|
|
switch name {
|
|
case "stopped":
|
|
s.p.initState = &stoppedState{p: s.p}
|
|
case "paused":
|
|
s.p.initState = &pausedState{p: s.p}
|
|
default:
|
|
return errors.Errorf("invalid state transition %q to %q", stateName(s), name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *runningState) Pause(ctx context.Context) error {
|
|
if err := s.p.runtime.Pause(ctx, s.p.id); err != nil {
|
|
return s.p.runtimeError(err, "OCI runtime pause failed")
|
|
}
|
|
|
|
return s.transition("paused")
|
|
}
|
|
|
|
func (s *runningState) Resume(ctx context.Context) error {
|
|
return errors.Errorf("cannot resume a running process")
|
|
}
|
|
|
|
func (s *runningState) Update(ctx context.Context, r *google_protobuf.Any) error {
|
|
return s.p.update(ctx, r)
|
|
}
|
|
|
|
func (s *runningState) Checkpoint(ctx context.Context, r *CheckpointConfig) error {
|
|
return s.p.checkpoint(ctx, r)
|
|
}
|
|
|
|
func (s *runningState) Resize(ws console.WinSize) error {
|
|
return s.p.resize(ws)
|
|
}
|
|
|
|
func (s *runningState) Start(ctx context.Context) error {
|
|
return errors.Errorf("cannot start a running process")
|
|
}
|
|
|
|
func (s *runningState) Delete(ctx context.Context) error {
|
|
return errors.Errorf("cannot delete a running process")
|
|
}
|
|
|
|
func (s *runningState) Kill(ctx context.Context, sig uint32, all bool) error {
|
|
return s.p.kill(ctx, sig, all)
|
|
}
|
|
|
|
func (s *runningState) SetExited(status int) {
|
|
s.p.setExited(status)
|
|
|
|
if err := s.transition("stopped"); err != nil {
|
|
panic(err)
|
|
}
|
|
}
|
|
|
|
func (s *runningState) Exec(ctx context.Context, path string, r *ExecConfig) (proc.Process, error) {
|
|
return s.p.exec(ctx, path, r)
|
|
}
|
|
|
|
type pausedState struct {
|
|
p *Init
|
|
}
|
|
|
|
func (s *pausedState) transition(name string) error {
|
|
switch name {
|
|
case "running":
|
|
s.p.initState = &runningState{p: s.p}
|
|
case "stopped":
|
|
s.p.initState = &stoppedState{p: s.p}
|
|
default:
|
|
return errors.Errorf("invalid state transition %q to %q", stateName(s), name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *pausedState) Pause(ctx context.Context) error {
|
|
return errors.Errorf("cannot pause a paused container")
|
|
}
|
|
|
|
func (s *pausedState) Resume(ctx context.Context) error {
|
|
if err := s.p.runtime.Resume(ctx, s.p.id); err != nil {
|
|
return s.p.runtimeError(err, "OCI runtime resume failed")
|
|
}
|
|
|
|
return s.transition("running")
|
|
}
|
|
|
|
func (s *pausedState) Update(ctx context.Context, r *google_protobuf.Any) error {
|
|
return s.p.update(ctx, r)
|
|
}
|
|
|
|
func (s *pausedState) Checkpoint(ctx context.Context, r *CheckpointConfig) error {
|
|
return s.p.checkpoint(ctx, r)
|
|
}
|
|
|
|
func (s *pausedState) Resize(ws console.WinSize) error {
|
|
return s.p.resize(ws)
|
|
}
|
|
|
|
func (s *pausedState) Start(ctx context.Context) error {
|
|
return errors.Errorf("cannot start a paused process")
|
|
}
|
|
|
|
func (s *pausedState) Delete(ctx context.Context) error {
|
|
return errors.Errorf("cannot delete a paused process")
|
|
}
|
|
|
|
func (s *pausedState) Kill(ctx context.Context, sig uint32, all bool) error {
|
|
return s.p.kill(ctx, sig, all)
|
|
}
|
|
|
|
func (s *pausedState) SetExited(status int) {
|
|
s.p.setExited(status)
|
|
|
|
if err := s.p.runtime.Resume(context.Background(), s.p.id); err != nil {
|
|
logrus.WithError(err).Error("resuming exited container from paused state")
|
|
}
|
|
|
|
if err := s.transition("stopped"); err != nil {
|
|
panic(err)
|
|
}
|
|
}
|
|
|
|
func (s *pausedState) Exec(ctx context.Context, path string, r *ExecConfig) (proc.Process, error) {
|
|
return nil, errors.Errorf("cannot exec in a paused state")
|
|
}
|
|
|
|
type stoppedState struct {
|
|
p *Init
|
|
}
|
|
|
|
func (s *stoppedState) transition(name string) error {
|
|
switch name {
|
|
case "deleted":
|
|
s.p.initState = &deletedState{}
|
|
default:
|
|
return errors.Errorf("invalid state transition %q to %q", stateName(s), name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *stoppedState) Pause(ctx context.Context) error {
|
|
return errors.Errorf("cannot pause a stopped container")
|
|
}
|
|
|
|
func (s *stoppedState) Resume(ctx context.Context) error {
|
|
return errors.Errorf("cannot resume a stopped container")
|
|
}
|
|
|
|
func (s *stoppedState) Update(ctx context.Context, r *google_protobuf.Any) error {
|
|
return errors.Errorf("cannot update a stopped container")
|
|
}
|
|
|
|
func (s *stoppedState) Checkpoint(ctx context.Context, r *CheckpointConfig) error {
|
|
return errors.Errorf("cannot checkpoint a stopped container")
|
|
}
|
|
|
|
func (s *stoppedState) Resize(ws console.WinSize) error {
|
|
return errors.Errorf("cannot resize a stopped container")
|
|
}
|
|
|
|
func (s *stoppedState) Start(ctx context.Context) error {
|
|
return errors.Errorf("cannot start a stopped process")
|
|
}
|
|
|
|
func (s *stoppedState) Delete(ctx context.Context) error {
|
|
if err := s.p.delete(ctx); err != nil {
|
|
return err
|
|
}
|
|
return s.transition("deleted")
|
|
}
|
|
|
|
func (s *stoppedState) Kill(ctx context.Context, sig uint32, all bool) error {
|
|
return errdefs.ToGRPCf(errdefs.ErrNotFound, "process %s not found", s.p.id)
|
|
}
|
|
|
|
func (s *stoppedState) SetExited(status int) {
|
|
// no op
|
|
}
|
|
|
|
func (s *stoppedState) Exec(ctx context.Context, path string, r *ExecConfig) (proc.Process, error) {
|
|
return nil, errors.Errorf("cannot exec in a stopped state")
|
|
}
|