package main
import (
"archive/tar"
"bytes"
"context"
"crypto/sha256"
"encoding/json"
"errors"
"fmt"
"io"
"math"
"math/rand"
"mime"
"mime/multipart"
"net/http"
"net/url"
"os"
"os/exec"
"path/filepath"
"slices"
"strconv"
"strings"
"sync"
"time"
"github.com/gorilla/mux"
"github.com/kballard/go-shellquote"
"gopkg.in/yaml.v2"
"github.com/lxc/incus/v6/client"
"github.com/lxc/incus/v6/internal/filter"
internalInstance "github.com/lxc/incus/v6/internal/instance"
internalIO "github.com/lxc/incus/v6/internal/io"
"github.com/lxc/incus/v6/internal/jmap"
"github.com/lxc/incus/v6/internal/server/auth"
"github.com/lxc/incus/v6/internal/server/cluster"
"github.com/lxc/incus/v6/internal/server/db"
dbCluster "github.com/lxc/incus/v6/internal/server/db/cluster"
"github.com/lxc/incus/v6/internal/server/db/operationtype"
"github.com/lxc/incus/v6/internal/server/instance"
"github.com/lxc/incus/v6/internal/server/instance/instancetype"
"github.com/lxc/incus/v6/internal/server/lifecycle"
"github.com/lxc/incus/v6/internal/server/node"
"github.com/lxc/incus/v6/internal/server/operations"
projectutils "github.com/lxc/incus/v6/internal/server/project"
"github.com/lxc/incus/v6/internal/server/request"
"github.com/lxc/incus/v6/internal/server/response"
"github.com/lxc/incus/v6/internal/server/state"
storagePools "github.com/lxc/incus/v6/internal/server/storage"
"github.com/lxc/incus/v6/internal/server/task"
localUtil "github.com/lxc/incus/v6/internal/server/util"
internalUtil "github.com/lxc/incus/v6/internal/util"
"github.com/lxc/incus/v6/internal/version"
"github.com/lxc/incus/v6/shared/api"
"github.com/lxc/incus/v6/shared/archive"
"github.com/lxc/incus/v6/shared/ioprogress"
"github.com/lxc/incus/v6/shared/logger"
"github.com/lxc/incus/v6/shared/osarch"
"github.com/lxc/incus/v6/shared/util"
)
var imagesCmd = APIEndpoint{
Path: "images",
Get: APIEndpointAction{Handler: imagesGet, AllowUntrusted: true},
Post: APIEndpointAction{Handler: imagesPost, AllowUntrusted: true},
}
var imageCmd = APIEndpoint{
Path: "images/{fingerprint}",
Delete: APIEndpointAction{Handler: imageDelete, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
Get: APIEndpointAction{Handler: imageGet, AllowUntrusted: true},
Patch: APIEndpointAction{Handler: imagePatch, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
Put: APIEndpointAction{Handler: imagePut, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
}
var imageExportCmd = APIEndpoint{
Path: "images/{fingerprint}/export",
Get: APIEndpointAction{Handler: imageExport, AllowUntrusted: true},
Post: APIEndpointAction{Handler: imageExportPost, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
}
var imageSecretCmd = APIEndpoint{
Path: "images/{fingerprint}/secret",
Post: APIEndpointAction{Handler: imageSecret, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
}
var imageRefreshCmd = APIEndpoint{
Path: "images/{fingerprint}/refresh",
Post: APIEndpointAction{Handler: imageRefresh, AccessHandler: allowPermission(auth.ObjectTypeImage, auth.EntitlementCanEdit, "fingerprint")},
}
var imageAliasesCmd = APIEndpoint{
Path: "images/aliases",
Get: APIEndpointAction{Handler: imageAliasesGet, AccessHandler: allowAuthenticated},
Post: APIEndpointAction{Handler: imageAliasesPost, AccessHandler: allowPermission(auth.ObjectTypeProject, auth.EntitlementCanCreateImageAliases)},
}
var imageAliasCmd = APIEndpoint{
Path: "images/aliases/{name:.*}",
Delete: APIEndpointAction{Handler: imageAliasDelete, AccessHandler: allowPermission(auth.ObjectTypeImageAlias, auth.EntitlementCanEdit, "name")},
Get: APIEndpointAction{Handler: imageAliasGet, AllowUntrusted: true},
Patch: APIEndpointAction{Handler: imageAliasPatch, AccessHandler: allowPermission(auth.ObjectTypeImageAlias, auth.EntitlementCanEdit, "name")},
Post: APIEndpointAction{Handler: imageAliasPost, AccessHandler: allowPermission(auth.ObjectTypeImageAlias, auth.EntitlementCanEdit, "name")},
Put: APIEndpointAction{Handler: imageAliasPut, AccessHandler: allowPermission(auth.ObjectTypeImageAlias, auth.EntitlementCanEdit, "name")},
}
We only want a single publish running at any one time.
The CPU and I/O load of publish is such that running multiple ones in
parallel takes longer than running them serially.
Additionally, publishing the same container or container snapshot
twice would lead to storage problem, not to mention a conflict at the
end for whichever finishes last.
*/
var imagePublishLock sync.Mutex
var imageTaskMu sync.Mutex
func compressFile(compress string, infile io.Reader, outfile io.Writer) error {
reproducible := []string{"gzip"}
var cmd *exec.Cmd
fields, err := shellquote.Split(compress)
if err != nil {
return err
}
if fields[0] == "squashfs" {
tempfile, err := os.CreateTemp("", "incus_compress_")
if err != nil {
return err
}
defer func() { _ = tempfile.Close() }()
defer func() { _ = os.Remove(tempfile.Name()) }()
args := []string{"tar2sqfs"}
if len(fields) > 1 {
args = append(args, fields[1:]...)
}
args = append(args, "--no-skip", "--force", "--compressor", "xz", tempfile.Name())
cmd = exec.Command(args[0], args[1:]...)
cmd.Stdin = infile
output, err := cmd.CombinedOutput()
if err != nil {
return fmt.Errorf("tar2sqfs: %v (%v)", err, strings.TrimSpace(string(output)))
}
_, err = tempfile.Seek(0, io.SeekStart)
if err != nil {
return err
}
_, err = io.Copy(outfile, tempfile)
if err != nil {
return err
}
} else {
args := []string{"-c"}
if len(fields) > 1 {
args = append(args, fields[1:]...)
}
if slices.Contains(reproducible, fields[0]) {
args = append(args, "-n")
}
cmd := exec.Command(fields[0], args...)
cmd.Stdin = infile
cmd.Stdout = outfile
err := cmd.Run()
if err != nil {
return err
}
}
return nil
}
* This function takes a container or snapshot from the local image server and
* exports it as an image.
*/
func imgPostInstanceInfo(ctx context.Context, s *state.State, r *http.Request, req api.ImagesPost, op *operations.Operation, builddir string, budget int64) (*api.Image, error) {
info := api.Image{}
info.Properties = map[string]string{}
projectName := request.ProjectParam(r)
name := req.Source.Name
ctype := req.Source.Type
if ctype == "" || name == "" {
return nil, fmt.Errorf("No source provided")
}
switch ctype {
case "snapshot":
if !internalInstance.IsSnapshot(name) {
return nil, fmt.Errorf("Not a snapshot")
}
case "container", "virtual-machine", "instance":
if internalInstance.IsSnapshot(name) {
return nil, fmt.Errorf("This is a snapshot")
}
default:
return nil, fmt.Errorf("Bad type")
}
info.Filename = req.Filename
switch req.Public {
case true:
info.Public = true
case false:
info.Public = false
}
c, err := instance.LoadByProjectAndName(s, projectName, name)
if err != nil {
return nil, err
}
info.Type = c.Type().String()
imageFile, err := os.CreateTemp(builddir, "incus_build_image_")
if err != nil {
return nil, err
}
defer func() { _ = os.Remove(imageFile.Name()) }()
totalSize := int64(0)
sumSize := func(path string, fi os.FileInfo, err error) error {
if err == nil {
totalSize += fi.Size()
}
return nil
}
err = filepath.Walk(c.RootfsPath(), sumSize)
if err != nil {
return nil, err
}
metadata := make(map[string]any)
imageProgressWriter := &ioprogress.ProgressWriter{
Tracker: &ioprogress.ProgressTracker{
Handler: func(value, speed int64) {
percent := int64(0)
var processed int64
if totalSize > 0 {
percent = value
processed = totalSize * (percent / 100.0)
} else {
processed = value
}
operations.SetProgressMetadata(metadata, "create_image_from_container_pack", "Image pack", percent, processed, speed)
_ = op.UpdateMetadata(metadata)
},
Length: totalSize,
},
}
sha256 := sha256.New()
var compress string
var writer io.Writer
if req.CompressionAlgorithm != "" {
compress = req.CompressionAlgorithm
} else {
var p *api.Project
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
project, err := dbCluster.GetProject(ctx, tx.Tx(), projectName)
if err != nil {
return err
}
p, err = project.ToAPI(ctx, tx.Tx())
return err
})
if err != nil {
return nil, err
}
if p.Config["images.compression_algorithm"] != "" {
compress = p.Config["images.compression_algorithm"]
} else {
compress = s.GlobalConfig.ImagesCompressionAlgorithm()
}
}
wg := sync.WaitGroup{}
var compressErr error
if compress != "none" {
wg.Add(1)
tarReader, tarWriter := io.Pipe()
imageProgressWriter.WriteCloser = tarWriter
writer = imageProgressWriter
compressWriter := io.MultiWriter(imageFile, sha256)
go func() {
defer wg.Done()
compressErr = compressFile(compress, tarReader, compressWriter)
if compressErr != nil {
_ = imageProgressWriter.Close()
}
}()
} else {
imageProgressWriter.WriteCloser = imageFile
writer = io.MultiWriter(imageProgressWriter, sha256)
}
tracker := &ioprogress.ProgressTracker{
Handler: func(value, speed int64) {
operations.SetProgressMetadata(metadata, "create_image_from_container_pack", "Exporting", value, 0, 0)
_ = op.UpdateMetadata(metadata)
},
}
var meta api.ImageMetadata
writer = internalIO.NewQuotaWriter(writer, budget)
meta, err = c.Export(writer, req.Properties, req.ExpiresAt, tracker)
if meta.ExpiryDate != 0 {
info.ExpiresAt = time.Unix(meta.ExpiryDate, 0)
}
_ = imageProgressWriter.Close()
wg.Wait()
_ = imageFile.Close()
if compressErr != nil {
return nil, compressErr
}
if err != nil {
return nil, err
}
fi, err := os.Stat(imageFile.Name())
if err != nil {
return nil, err
}
info.Size = fi.Size()
info.Fingerprint = fmt.Sprintf("%x", sha256.Sum(nil))
info.CreatedAt = time.Now().UTC()
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
_, _, err = tx.GetImage(ctx, info.Fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if !response.IsNotFoundError(err) {
if err != nil {
return nil, err
}
return &info, fmt.Errorf("The image already exists: %s", info.Fingerprint)
}
finalName := internalUtil.VarPath("images", info.Fingerprint)
err = internalUtil.FileMove(imageFile.Name(), finalName)
if err != nil {
return nil, err
}
info.Architecture, _ = osarch.ArchitectureName(c.Architecture())
info.Properties = meta.Properties
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.CreateImage(ctx, c.Project().Name, info.Fingerprint, info.Filename, info.Size, info.Public, info.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, info.Properties, info.Type, nil)
})
if err != nil {
return nil, err
}
return &info, nil
}
func imgPostRemoteInfo(ctx context.Context, s *state.State, r *http.Request, req api.ImagesPost, op *operations.Operation, project string, budget int64) (*api.Image, error) {
var err error
var hash string
if req.Source.Fingerprint != "" {
hash = req.Source.Fingerprint
} else if req.Source.Alias != "" {
hash = req.Source.Alias
} else {
return nil, fmt.Errorf("must specify one of alias or fingerprint for init from image")
}
info, err := ImageDownload(ctx, r, s, op, &ImageDownloadArgs{
Server: req.Source.Server,
Protocol: req.Source.Protocol,
Certificate: req.Source.Certificate,
Secret: req.Source.Secret,
Alias: hash,
Type: req.Source.ImageType,
AutoUpdate: req.AutoUpdate,
Public: req.Public,
ProjectName: project,
Budget: budget,
SourceProjectName: req.Source.Project,
})
if err != nil {
return nil, err
}
if isClusterNotification(r) && req.Source.Server == "" {
return info, nil
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var id int
id, info, err = tx.GetImage(ctx, info.Fingerprint, dbCluster.ImageFilter{Project: &project})
if err != nil {
return err
}
for k, v := range req.Properties {
info.Properties[k] = v
}
if req.Profiles == nil {
req.Profiles = []string{api.ProjectDefaultName}
}
profileIds := make([]int64, len(req.Profiles))
for i, profile := range req.Profiles {
profileID, _, err := tx.GetProfile(ctx, project, profile)
if response.IsNotFoundError(err) {
return fmt.Errorf("Profile '%s' doesn't exist", profile)
} else if err != nil {
return err
}
profileIds[i] = profileID
}
if req.Public || req.AutoUpdate || req.Filename != "" || len(req.Properties) > 0 || len(req.Profiles) > 0 {
err := tx.UpdateImage(ctx, id, req.Filename, info.Size, req.Public, req.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, info.Properties, project, profileIds)
if err != nil {
return err
}
}
return nil
})
if err != nil {
return nil, err
}
return info, nil
}
func imgPostURLInfo(ctx context.Context, s *state.State, r *http.Request, req api.ImagesPost, op *operations.Operation, project string, budget int64) (*api.Image, error) {
var err error
if req.Source.URL == "" {
return nil, fmt.Errorf("Missing URL")
}
myhttp, err := localUtil.HTTPClient("", s.Proxy)
if err != nil {
return nil, err
}
head, err := http.NewRequest("HEAD", req.Source.URL, nil)
if err != nil {
return nil, err
}
architectures := []string{}
for _, architecture := range s.OS.Architectures {
architectureName, err := osarch.ArchitectureName(architecture)
if err != nil {
return nil, err
}
architectures = append(architectures, architectureName)
}
head.Header.Set("User-Agent", version.UserAgent)
head.Header.Set("Incus-Server-Architectures", strings.Join(architectures, ", "))
head.Header.Set("Incus-Server-Version", version.Version)
raw, err := myhttp.Do(head)
if err != nil {
return nil, err
}
hash := raw.Header.Get("Incus-Image-Hash")
if hash == "" {
return nil, fmt.Errorf("Missing Incus-Image-Hash header")
}
url := raw.Header.Get("Incus-Image-URL")
if url == "" {
return nil, fmt.Errorf("Missing Incus-Image-URL header")
}
info, err := ImageDownload(ctx, r, s, op, &ImageDownloadArgs{
Server: url,
Protocol: "direct",
Alias: hash,
AutoUpdate: req.AutoUpdate,
Public: req.Public,
ProjectName: project,
Budget: budget,
})
if err != nil {
return nil, err
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var id int
id, info, err = tx.GetImage(ctx, info.Fingerprint, dbCluster.ImageFilter{Project: &project})
if err != nil {
return err
}
for k, v := range req.Properties {
info.Properties[k] = v
}
if req.Public || req.AutoUpdate || req.Filename != "" || len(req.Properties) > 0 {
err := tx.UpdateImage(ctx, id, req.Filename, info.Size, req.Public, req.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, info.Properties, "", nil)
if err != nil {
return err
}
}
return nil
})
if err != nil {
return nil, err
}
return info, nil
}
func getImgPostInfo(ctx context.Context, s *state.State, r *http.Request, builddir string, project string, post *os.File, metadata map[string]any) (*api.Image, error) {
info := api.Image{}
var imageMeta *api.ImageMetadata
l := logger.AddContext(logger.Ctx{"function": "getImgPostInfo"})
info.Public = util.IsTrue(r.Header.Get("X-Incus-public"))
propHeaders := r.Header[http.CanonicalHeaderKey("X-Incus-properties")]
profilesHeaders := r.Header.Get("X-Incus-profiles")
ctype, ctypeParams, err := mime.ParseMediaType(r.Header.Get("Content-Type"))
if err != nil {
ctype = "application/octet-stream"
}
sha256 := sha256.New()
var size int64
if ctype == "multipart/form-data" {
imageTarf, err := os.CreateTemp(builddir, "incus_tar_")
if err != nil {
return nil, err
}
defer func() { _ = os.Remove(imageTarf.Name()) }()
_, err = post.Seek(0, io.SeekStart)
if err != nil {
return nil, err
}
mr := multipart.NewReader(post, ctypeParams["boundary"])
part, err := mr.NextPart()
if err != nil {
return nil, err
}
if part.FormName() != "metadata" {
return nil, fmt.Errorf("Invalid multipart image")
}
size, err = io.Copy(io.MultiWriter(imageTarf, sha256), part)
info.Size += size
_ = imageTarf.Close()
if err != nil {
l.Error("Failed to copy the image tarfile", logger.Ctx{"err": err})
return nil, err
}
part, err = mr.NextPart()
if err != nil {
l.Error("Failed to get the next part", logger.Ctx{"err": err})
return nil, err
}
if part.FormName() == "rootfs" {
info.Type = instancetype.Container.String()
} else if part.FormName() == "rootfs.img" {
info.Type = instancetype.VM.String()
} else {
l.Error("Invalid multipart image")
return nil, fmt.Errorf("Invalid multipart image")
}
rootfsTarf, err := os.CreateTemp(builddir, "incus_tar_")
if err != nil {
return nil, err
}
defer func() { _ = os.Remove(rootfsTarf.Name()) }()
size, err = io.Copy(io.MultiWriter(rootfsTarf, sha256), part)
info.Size += size
_ = rootfsTarf.Close()
if err != nil {
l.Error("Failed to copy the rootfs tarfile", logger.Ctx{"err": err})
return nil, err
}
info.Filename = part.FileName()
info.Fingerprint = fmt.Sprintf("%x", sha256.Sum(nil))
expectedFingerprint := r.Header.Get("X-Incus-fingerprint")
if expectedFingerprint != "" && info.Fingerprint != expectedFingerprint {
err = fmt.Errorf("fingerprints don't match, got %s expected %s", info.Fingerprint, expectedFingerprint)
return nil, err
}
imageMeta, _, err = getImageMetadata(imageTarf.Name())
if err != nil {
l.Error("Failed to get image metadata", logger.Ctx{"err": err})
return nil, err
}
imgfname := internalUtil.VarPath("images", info.Fingerprint)
err = internalUtil.FileMove(imageTarf.Name(), imgfname)
if err != nil {
l.Error("Failed to move the image tarfile", logger.Ctx{
"err": err,
"source": imageTarf.Name(),
"dest": imgfname})
return nil, err
}
rootfsfname := internalUtil.VarPath("images", info.Fingerprint+".rootfs")
err = internalUtil.FileMove(rootfsTarf.Name(), rootfsfname)
if err != nil {
l.Error("Failed to move the rootfs tarfile", logger.Ctx{
"err": err,
"source": rootfsTarf.Name(),
"dest": imgfname})
return nil, err
}
} else {
_, err = post.Seek(0, io.SeekStart)
if err != nil {
return nil, err
}
size, err = io.Copy(sha256, post)
if err != nil {
l.Error("Failed to copy the tarfile", logger.Ctx{"err": err})
return nil, err
}
info.Size = size
info.Filename = r.Header.Get("X-Incus-filename")
info.Fingerprint = fmt.Sprintf("%x", sha256.Sum(nil))
expectedFingerprint := r.Header.Get("X-Incus-fingerprint")
if expectedFingerprint != "" && info.Fingerprint != expectedFingerprint {
l.Error("Fingerprints don't match", logger.Ctx{
"got": info.Fingerprint,
"expected": expectedFingerprint})
err = fmt.Errorf("fingerprints don't match, got %s expected %s", info.Fingerprint, expectedFingerprint)
return nil, err
}
var imageType string
imageMeta, imageType, err = getImageMetadata(post.Name())
if err != nil {
l.Error("Failed to get image metadata", logger.Ctx{"err": err})
return nil, err
}
info.Type = imageType
imgfname := internalUtil.VarPath("images", info.Fingerprint)
err = internalUtil.FileMove(post.Name(), imgfname)
if err != nil {
l.Error("Failed to move the tarfile", logger.Ctx{
"err": err,
"source": post.Name(),
"dest": imgfname})
return nil, err
}
}
info.Architecture = imageMeta.Architecture
if imageMeta.CreationDate > 0 {
info.CreatedAt = time.Unix(imageMeta.CreationDate, 0)
}
expiresAt, ok := metadata["expires_at"]
if ok {
info.ExpiresAt = expiresAt.(time.Time)
} else if imageMeta.ExpiryDate > 0 {
info.ExpiresAt = time.Unix(imageMeta.ExpiryDate, 0)
}
properties, ok := metadata["properties"]
if ok {
info.Properties = properties.(map[string]string)
} else {
info.Properties = imageMeta.Properties
}
if len(propHeaders) > 0 {
for _, ph := range propHeaders {
p, _ := url.ParseQuery(ph)
for pkey, pval := range p {
info.Properties[pkey] = pval[0]
}
}
}
var profileIds []int64
if len(profilesHeaders) > 0 {
p, _ := url.ParseQuery(profilesHeaders)
profileIds = make([]int64, len(p["profile"]))
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
for i, val := range p["profile"] {
profileID, _, err := tx.GetProfile(ctx, project, val)
if response.IsNotFoundError(err) {
return fmt.Errorf("Profile '%s' doesn't exist", val)
} else if err != nil {
return err
}
profileIds[i] = profileID
}
return nil
})
if err != nil {
return nil, err
}
}
var exists bool
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
exists, err = tx.ImageExists(ctx, project, info.Fingerprint)
return err
})
if err != nil {
return nil, err
}
if exists {
if isClusterNotification(r) {
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.AddImageToLocalNode(ctx, project, info.Fingerprint)
})
if err != nil {
return nil, err
}
} else {
return &info, fmt.Errorf("Image with same fingerprint already exists")
}
} else {
public, ok := metadata["public"]
if ok {
info.Public = public.(bool)
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.CreateImage(ctx, project, info.Fingerprint, info.Filename, info.Size, info.Public, info.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, info.Properties, info.Type, profileIds)
})
if err != nil {
return nil, err
}
}
return &info, nil
}
func imageCreateInPool(s *state.State, info *api.Image, storagePool string) error {
if storagePool == "" {
return fmt.Errorf("No storage pool specified")
}
pool, err := storagePools.LoadByName(s, storagePool)
if err != nil {
return err
}
err = pool.EnsureImage(info.Fingerprint, nil)
if err != nil {
return err
}
return nil
}
func imagesPost(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
var userCanCreateImages bool
err := s.Authorizer.CheckPermission(r.Context(), r, auth.ObjectProject(projectName), auth.EntitlementCanCreateImages)
if err == nil {
userCanCreateImages = true
} else if !api.StatusErrorCheck(err, http.StatusForbidden) {
return response.SmartError(err)
}
trusted := d.checkTrustedClient(r) == nil && userCanCreateImages
secret := r.Header.Get("X-Incus-secret")
fingerprint := r.Header.Get("X-Incus-fingerprint")
var imageMetadata map[string]any
if !trusted && (secret == "" || fingerprint == "") {
return response.Forbidden(nil)
} else {
op, err := imageValidSecret(s, r, projectName, fingerprint, secret)
if err != nil {
return response.SmartError(err)
}
if op != nil {
imageMetadata = op.Metadata
} else if !trusted {
return response.Forbidden(nil)
}
}
instanceType, err := urlInstanceTypeDetect(r)
if err != nil {
return response.SmartError(err)
}
builddir, err := os.MkdirTemp(internalUtil.VarPath("images"), "incus_build_")
if err != nil {
return response.InternalError(err)
}
cleanup := func(path string, fd *os.File) {
if fd != nil {
_ = fd.Close()
}
err := os.RemoveAll(path)
if err != nil {
logger.Debugf("Error deleting temporary directory \"%s\": %s", path, err)
}
}
post, err := os.CreateTemp(builddir, "incus_post_")
if err != nil {
cleanup(builddir, nil)
return response.InternalError(err)
}
var budget int64
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
budget, err = projectutils.GetImageSpaceBudget(tx, projectName)
return err
})
if err != nil {
return response.SmartError(err)
}
_, err = io.Copy(internalIO.NewQuotaWriter(post, budget), r.Body)
if err != nil {
logger.Errorf("Store image POST data to disk: %v", err)
cleanup(builddir, post)
return response.InternalError(err)
}
_, err = post.Seek(0, io.SeekStart)
if err != nil {
return response.InternalError(err)
}
decoder := json.NewDecoder(post)
imageUpload := false
req := api.ImagesPost{}
err = decoder.Decode(&req)
if err != nil {
if r.Header.Get("Content-Type") == "application/json" {
cleanup(builddir, post)
return response.BadRequest(err)
}
imageUpload = true
}
if !imageUpload && req.Source.Mode == "push" {
cleanup(builddir, post)
metadata := map[string]any{
"aliases": req.Aliases,
"expires_at": req.ExpiresAt,
"properties": req.Properties,
"public": req.Public,
}
return createTokenResponse(s, r, projectName, req.Source.Fingerprint, metadata)
}
if !imageUpload && !slices.Contains([]string{"container", "instance", "virtual-machine", "snapshot", "image", "url"}, req.Source.Type) {
cleanup(builddir, post)
return response.InternalError(fmt.Errorf("Invalid images JSON"))
}
if !imageUpload && slices.Contains([]string{"container", "instance", "virtual-machine", "snapshot"}, req.Source.Type) {
name := req.Source.Name
if name != "" {
_, err = post.Seek(0, io.SeekStart)
if err != nil {
return response.InternalError(err)
}
r.Body = post
resp, err := forwardedResponseIfInstanceIsRemote(s, r, projectName, name, instanceType)
if err != nil {
cleanup(builddir, post)
return response.SmartError(err)
}
if resp != nil {
cleanup(builddir, nil)
return resp
}
}
}
run := func(op *operations.Operation) error {
var err error
var info *api.Image
defer cleanup(builddir, post)
if imageUpload {
info, err = getImgPostInfo(context.TODO(), s, r, builddir, projectName, post, imageMetadata)
} else {
if req.Source.Type == "image" {
info, err = imgPostRemoteInfo(context.TODO(), s, r, req, op, projectName, budget)
} else if req.Source.Type == "url" {
info, err = imgPostURLInfo(context.TODO(), s, r, req, op, projectName, budget)
} else {
imagePublishLock.Lock()
info, err = imgPostInstanceInfo(context.TODO(), s, r, req, op, builddir, budget)
imagePublishLock.Unlock()
}
}
if info != nil {
metadata := make(map[string]string)
metadata["fingerprint"] = info.Fingerprint
metadata["size"] = strconv.FormatInt(info.Size, 10)
secret, ok := op.Metadata()["secret"]
if ok {
metadata["secret"] = secret.(string)
}
_ = op.UpdateMetadata(metadata)
}
if err != nil {
return err
}
if isClusterNotification(r) {
return nil
}
aliases, ok := imageMetadata["aliases"]
if ok {
req.Aliases = aliases.([]api.ImageAlias)
}
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
imgID, _, err := tx.GetImageByFingerprintPrefix(ctx, info.Fingerprint, dbCluster.ImageFilter{Project: &projectName})
if err != nil {
return fmt.Errorf("Fetch image %q: %w", info.Fingerprint, err)
}
for _, alias := range req.Aliases {
_, _, err := tx.GetImageAlias(ctx, projectName, alias.Name, true)
if !response.IsNotFoundError(err) {
if err != nil {
return fmt.Errorf("Fetch image alias %q: %w", alias.Name, err)
}
return fmt.Errorf("Alias already exists: %s", alias.Name)
}
err = tx.CreateImageAlias(ctx, projectName, alias.Name, imgID, alias.Description)
if err != nil {
return fmt.Errorf("Add new image alias to the database: %w", err)
}
err = s.Authorizer.AddImageAlias(r.Context(), projectName, alias.Name)
if err != nil {
logger.Error("Failed to add image alias to authorizer", logger.Ctx{"name": alias.Name, "project": projectName, "error": err})
}
s.Events.SendLifecycle(projectName, lifecycle.ImageAliasCreated.Event(alias.Name, projectName, op.Requestor(), logger.Ctx{"target": info.Fingerprint}))
}
return nil
})
if err != nil {
return err
}
err = imageSyncBetweenNodes(context.TODO(), s, r, projectName, info.Fingerprint)
if err != nil {
return fmt.Errorf("Failed syncing image between servers: %w", err)
}
err = s.Authorizer.AddImage(s.ShutdownCtx, projectName, info.Fingerprint)
if err != nil {
logger.Error("Failed to add image to authorizer", logger.Ctx{"fingerprint": info.Fingerprint, "project": projectName, "error": err})
}
s.Events.SendLifecycle(projectName, lifecycle.ImageCreated.Event(info.Fingerprint, projectName, op.Requestor(), logger.Ctx{"type": info.Type}))
return nil
}
var metadata any
if imageUpload && imageMetadata != nil {
secret, _ := internalUtil.RandomHexString(32)
if secret != "" {
metadata = map[string]string{
"secret": secret,
}
}
}
op, err := operations.OperationCreate(s, projectName, operations.OperationClassTask, operationtype.ImageDownload, nil, metadata, run, nil, nil, r)
if err != nil {
cleanup(builddir, post)
return response.InternalError(err)
}
return operations.OperationResponse(op)
}
func getImageMetadata(fname string) (*api.ImageMetadata, string, error) {
var tr *tar.Reader
var result api.ImageMetadata
r, err := os.Open(fname)
if err != nil {
return nil, "unknown", err
}
defer func() { _ = r.Close() }()
_, algo, unpacker, err := archive.DetectCompressionFile(r)
if err != nil {
return nil, "unknown", err
}
_, err = r.Seek(0, io.SeekStart)
if err != nil {
return nil, "", err
}
if unpacker == nil {
return nil, "unknown", fmt.Errorf("Unsupported backup compression")
}
if len(unpacker) > 0 {
if algo == ".squashfs" {
unpacker = append(unpacker, fname)
}
cmd := exec.Command(unpacker[0], unpacker[1:]...)
cmd.Stdin = r
stdout, err := cmd.StdoutPipe()
if err != nil {
return nil, "unknown", err
}
defer func() { _ = stdout.Close() }()
err = cmd.Start()
if err != nil {
return nil, "unknown", err
}
defer func() { _ = cmd.Wait() }()
defer func() { _ = stdout.Close() }()
tr = tar.NewReader(stdout)
} else {
tr = tar.NewReader(r)
}
hasMeta := false
hasRoot := false
imageType := "unknown"
for {
hdr, err := tr.Next()
if err == io.EOF {
break
}
if err != nil {
return nil, "unknown", err
}
if hdr.Name == "metadata.yaml" || hdr.Name == "./metadata.yaml" {
err = yaml.NewDecoder(tr).Decode(&result)
if err != nil {
return nil, "unknown", err
}
hasMeta = true
}
if strings.HasPrefix(hdr.Name, "rootfs/") || strings.HasPrefix(hdr.Name, "./rootfs/") {
hasRoot = true
imageType = instancetype.Container.String()
}
if hdr.Name == "rootfs.img" || hdr.Name == "./rootfs.img" {
hasRoot = true
imageType = instancetype.VM.String()
}
if hasMeta && hasRoot {
break
}
}
if !hasMeta {
return nil, "unknown", fmt.Errorf("Metadata tarball is missing metadata.yaml")
}
_, err = osarch.ArchitectureId(result.Architecture)
if err != nil {
return nil, "unknown", err
}
if result.CreationDate == 0 {
return nil, "unknown", fmt.Errorf("Missing creation date")
}
return &result, imageType, nil
}
func doImagesGet(ctx context.Context, tx *db.ClusterTx, recursion bool, projectName string, public bool, clauses *filter.ClauseSet, hasPermission auth.PermissionChecker, allProjects bool) (any, error) {
mustLoadObjects := recursion || (clauses != nil && len(clauses.Clauses) > 0)
imagesProjectsMap := map[string][]string{}
if allProjects {
var err error
imagesProjectsMap, err = tx.GetImages(ctx)
if err != nil {
return nil, err
}
} else {
fingerprints, err := tx.GetImagesFingerprints(ctx, projectName, public)
if err != nil {
return nil, err
}
for _, fp := range fingerprints {
imagesProjectsMap[fp] = []string{projectName}
}
}
var resultString []string
var resultMap []*api.Image
if recursion {
resultMap = make([]*api.Image, 0, len(imagesProjectsMap))
} else {
resultString = make([]string, 0, len(imagesProjectsMap))
}
for fingerprint, projects := range imagesProjectsMap {
for _, curProjectName := range projects {
image, err := doImageGet(ctx, tx, curProjectName, fingerprint, public)
if err != nil {
continue
}
if !image.Public && !hasPermission(auth.ObjectImage(curProjectName, fingerprint)) {
continue
}
if !mustLoadObjects {
resultString = append(resultString, api.NewURL().Path(version.APIVersion, "images", fingerprint).String())
} else {
if clauses != nil && len(clauses.Clauses) > 0 {
match, err := filter.Match(*image, *clauses)
if err != nil {
return nil, err
}
if !match {
continue
}
}
if recursion {
resultMap = append(resultMap, image)
} else {
resultString = append(resultString, api.NewURL().Path(version.APIVersion, "images", image.Fingerprint).String())
}
}
}
}
if recursion {
return resultMap, nil
}
return resultString, nil
}
func imagesGet(d *Daemon, r *http.Request) response.Response {
projectName := request.ProjectParam(r)
allProjects := util.IsTrue(r.FormValue("all-projects"))
filterStr := r.FormValue("filter")
if allProjects && projectName != api.ProjectDefaultName {
return response.BadRequest(fmt.Errorf("Cannot specify a project when requesting all projects"))
}
s := d.State()
hasPermission, authorizationErr := s.Authorizer.GetPermissionChecker(r.Context(), r, auth.EntitlementCanView, auth.ObjectTypeImage)
if authorizationErr != nil && !api.StatusErrorCheck(authorizationErr, http.StatusForbidden) {
return response.SmartError(authorizationErr)
}
public := d.checkTrustedClient(r) != nil || authorizationErr != nil
clauses, err := filter.Parse(filterStr, filter.QueryOperatorSet())
if err != nil {
return response.SmartError(fmt.Errorf("Invalid filter: %w", err))
}
var result any
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
result, err = doImagesGet(ctx, tx, localUtil.IsRecursionRequest(r), projectName, public, clauses, hasPermission, allProjects)
if err != nil {
return err
}
return nil
})
if err != nil {
return response.SmartError(err)
}
return response.SyncResponse(true, result)
}
func autoUpdateImagesTask(d *Daemon) (task.Func, task.Schedule) {
f := func(ctx context.Context) {
s := d.State()
opRun := func(op *operations.Operation) error {
return autoUpdateImages(ctx, s)
}
op, err := operations.OperationCreate(s, "", operations.OperationClassTask, operationtype.ImagesUpdate, nil, nil, opRun, nil, nil, nil)
if err != nil {
logger.Error("Failed creating image update operation", logger.Ctx{"err": err})
return
}
logger.Debug("Acquiring image task lock")
imageTaskMu.Lock()
defer imageTaskMu.Unlock()
logger.Debug("Acquired image task lock")
logger.Info("Updating images")
err = op.Start()
if err != nil {
logger.Error("Failed starting image update operation", logger.Ctx{"err": err})
return
}
err = op.Wait(ctx)
if err != nil {
logger.Error("Failed updating images", logger.Ctx{"err": err})
return
}
logger.Info("Done updating images")
}
return f, task.Hourly()
}
func autoUpdateImages(ctx context.Context, s *state.State) error {
imageMap := make(map[string][]dbCluster.Image)
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
autoUpdate := true
images, err := dbCluster.GetImages(ctx, tx.Tx(), dbCluster.ImageFilter{AutoUpdate: &autoUpdate})
if err != nil {
return err
}
for _, image := range images {
imageMap[image.Fingerprint] = append(imageMap[image.Fingerprint], image)
}
return nil
})
if err != nil {
return fmt.Errorf("Unable to retrieve image fingerprints: %w", err)
}
for fingerprint, images := range imageMap {
skipFingerprint := false
var nodes []string
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
nodes, err = tx.GetNodesWithImageAndAutoUpdate(ctx, fingerprint, true)
return err
})
if err != nil {
logger.Error("Error getting cluster members for image auto-update", logger.Ctx{"fingerprint": fingerprint, "err": err})
continue
}
if len(nodes) > 1 {
var nodeIDs []int64
for _, node := range nodes {
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
nodeInfo, err := tx.GetNodeByAddress(ctx, node)
if err != nil {
return err
}
nodeIDs = append(nodeIDs, nodeInfo.ID)
return nil
})
if err != nil {
logger.Error("Unable to retrieve cluster member information for image update", logger.Ctx{"err": err})
skipFingerprint = true
break
}
}
if skipFingerprint {
continue
}
if len(nodeIDs) > 1 {
selectedNode, err := localUtil.GetStableRandomInt64FromList(int64(len(images)), nodeIDs)
if err != nil {
logger.Error("Failed to select cluster member for image update", logger.Ctx{"err": err})
continue
}
if s.DB.Cluster.GetNodeID() != selectedNode {
continue
}
}
}
var deleteIDs []int
var newImage *api.Image
for _, image := range images {
filter := dbCluster.ImageFilter{Project: &image.Project}
if image.Public {
filter.Public = &image.Public
}
var imageInfo *api.Image
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
_, imageInfo, err = tx.GetImage(ctx, image.Fingerprint, filter)
return err
})
if err != nil {
logger.Error("Failed to get image", logger.Ctx{"err": err, "project": image.Project, "fingerprint": image.Fingerprint})
continue
}
newInfo, err := autoUpdateImage(ctx, s, nil, image.ID, imageInfo, image.Project, false)
if err != nil {
logger.Error("Failed to update image", logger.Ctx{"err": err, "project": image.Project, "fingerprint": image.Fingerprint})
if err == context.Canceled {
return nil
}
} else {
deleteIDs = append(deleteIDs, image.ID)
}
if newImage == nil {
newImage = newInfo
}
}
if newImage != nil {
if len(nodes) > 1 {
err := distributeImage(ctx, s, nodes, fingerprint, newImage)
if err != nil {
logger.Error("Failed to distribute new image", logger.Ctx{"err": err, "fingerprint": newImage.Fingerprint})
if err == context.Canceled {
return nil
}
}
}
_ = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
for _, ID := range deleteIDs {
err := tx.DeleteImage(ctx, ID)
if err != nil {
logger.Error("Error deleting old image from database", logger.Ctx{"err": err, "fingerprint": fingerprint, "ID": ID})
}
}
return nil
})
}
}
return nil
}
func distributeImage(ctx context.Context, s *state.State, nodes []string, oldFingerprint string, newImage *api.Image) error {
var imageVolumes []string
err := s.DB.Node.Transaction(ctx, func(ctx context.Context, tx *db.NodeTx) error {
config, err := node.ConfigLoad(ctx, tx)
if err != nil {
return err
}
vol := config.StorageImagesVolume()
if vol != "" {
fields := strings.Split(vol, "/")
var pool *api.StoragePool
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
_, pool, _, err = tx.GetStoragePool(ctx, fields[0])
return err
})
if err != nil {
return fmt.Errorf("Failed to get storage pool info: %w", err)
}
if slices.Contains(db.StorageRemoteDriverNames(), pool.Driver) {
imageVolumes = append(imageVolumes, vol)
}
}
return nil
})
if err != nil {
logger.Error("Failed to load config", logger.Ctx{"err": err})
}
localClusterAddress := s.LocalConfig.ClusterAddress()
var poolIDs []int64
var poolNames []string
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
poolIDs, err = tx.GetPoolsWithImage(ctx, newImage.Fingerprint)
if err != nil {
logger.Error("Error getting image storage pools", logger.Ctx{"err": err, "fingerprint": oldFingerprint})
return err
}
poolNames, err = tx.GetPoolNamesFromIDs(ctx, poolIDs)
if err != nil {
logger.Error("Error getting image storage pools", logger.Ctx{"err": err, "fingerprint": oldFingerprint})
return err
}
return nil
})
if err != nil {
return err
}
for _, nodeAddress := range nodes {
if nodeAddress == localClusterAddress {
continue
}
var nodeInfo db.NodeInfo
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
nodeInfo, err = tx.GetNodeByAddress(ctx, nodeAddress)
return err
})
if err != nil {
return fmt.Errorf("Failed to retrieve information about cluster member with address %q: %w", nodeAddress, err)
}
client, err := cluster.Connect(nodeAddress, s.Endpoints.NetworkCert(), s.ServerCert(), nil, true)
if err != nil {
return fmt.Errorf("Failed to connect to %q for image synchronization: %w", nodeAddress, err)
}
client = client.UseTarget(nodeInfo.Name)
resp, _, err := client.GetServer()
if err != nil {
logger.Error("Failed to retrieve information about cluster member", logger.Ctx{"err": err, "remote": nodeAddress})
} else {
vol := resp.Config["storage.images_volume"]
skipDistribution := false
if vol != "" {
for _, imageVolume := range imageVolumes {
if imageVolume == vol {
skipDistribution = true
break
}
}
if skipDistribution {
continue
}
fields := strings.Split(vol, "/")
pool, _, err := client.GetStoragePool(fields[0])
if err != nil {
logger.Error("Failed to get storage pool info", logger.Ctx{"err": err, "pool": fields[0]})
} else {
if slices.Contains(db.StorageRemoteDriverNames(), pool.Driver) {
imageVolumes = append(imageVolumes, vol)
}
}
}
}
createArgs := &incus.ImageCreateArgs{}
imageMetaPath := internalUtil.VarPath("images", newImage.Fingerprint)
imageRootfsPath := internalUtil.VarPath("images", newImage.Fingerprint+".rootfs")
metaFile, err := os.Open(imageMetaPath)
if err != nil {
return err
}
defer func() { _ = metaFile.Close() }()
createArgs.MetaFile = metaFile
createArgs.MetaName = filepath.Base(imageMetaPath)
createArgs.Type = newImage.Type
if util.PathExists(imageRootfsPath) {
rootfsFile, err := os.Open(imageRootfsPath)
if err != nil {
return err
}
defer func() { _ = rootfsFile.Close() }()
createArgs.RootfsFile = rootfsFile
createArgs.RootfsName = filepath.Base(imageRootfsPath)
}
image := api.ImagesPost{}
image.Filename = createArgs.MetaName
op, err := client.CreateImage(image, createArgs)
if err != nil {
return err
}
select {
case <-ctx.Done():
_ = op.Cancel()
return ctx.Err()
default:
}
err = op.Wait()
if err != nil {
return err
}
for _, poolName := range poolNames {
if poolName == "" {
continue
}
req := internalImageOptimizePost{
Image: *newImage,
Pool: poolName,
}
_, _, err = client.RawQuery("POST", "/internal/image-optimize", req, "")
if err != nil {
logger.Error("Failed creating new image in storage pool", logger.Ctx{"err": err, "remote": nodeAddress, "pool": poolName, "fingerprint": newImage.Fingerprint})
}
err = client.DeleteStoragePoolVolume(poolName, "image", oldFingerprint)
if err != nil {
logger.Error("Failed deleting old image from storage pool", logger.Ctx{"err": err, "remote": nodeAddress, "pool": poolName, "fingerprint": oldFingerprint})
}
}
}
return nil
}
func autoUpdateImage(ctx context.Context, s *state.State, op *operations.Operation, id int, info *api.Image, projectName string, manual bool) (*api.Image, error) {
fingerprint := info.Fingerprint
var source api.ImageSource
if !manual {
var interval int64
var project *api.Project
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
p, err := dbCluster.GetProject(ctx, tx.Tx(), projectName)
if err != nil {
return err
}
project, err = p.ToAPI(ctx, tx.Tx())
return err
})
if err != nil {
return nil, err
}
if project.Config["images.auto_update_interval"] != "" {
interval, err = strconv.ParseInt(project.Config["images.auto_update_interval"], 10, 64)
if err != nil {
return nil, fmt.Errorf("Unable to fetch project configuration: %w", err)
}
} else {
interval = s.GlobalConfig.ImagesAutoUpdateIntervalHours()
}
if interval <= 0 {
return nil, nil
}
now := time.Now()
elapsedHours := int64(math.Round(now.Sub(s.StartTime).Hours()))
if elapsedHours%interval != 0 {
return nil, nil
}
}
var poolNames []string
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
_, source, err = tx.GetImageSource(ctx, id)
if err != nil {
logger.Error("Error getting source image", logger.Ctx{"err": err, "fingerprint": fingerprint})
return err
}
poolIDs, err := tx.GetPoolsWithImage(ctx, fingerprint)
if err != nil {
logger.Error("Error getting image pools", logger.Ctx{"err": err, "fingerprint": fingerprint})
return err
}
poolNames, err = tx.GetPoolNamesFromIDs(ctx, poolIDs)
if err != nil {
logger.Error("Error getting image pools", logger.Ctx{"err": err, "fingerprint": fingerprint})
return err
}
return nil
})
if err != nil {
return nil, err
}
if len(poolNames) == 0 {
poolNames = append(poolNames, "")
}
logger.Debug("Processing image", logger.Ctx{"fingerprint": fingerprint, "server": source.Server, "protocol": source.Protocol, "alias": source.Alias})
setRefreshResult := func(result bool) {
if op == nil {
return
}
metadata := map[string]any{"refreshed": result}
_ = op.UpdateMetadata(metadata)
if result {
s.Events.SendLifecycle(projectName, lifecycle.ImageRefreshed.Event(fingerprint, projectName, op.Requestor(), nil))
}
}
hash := fingerprint
var newInfo *api.Image
for _, poolName := range poolNames {
select {
case <-ctx.Done():
return nil, ctx.Err()
default:
}
newInfo, err = ImageDownload(ctx, nil, s, op, &ImageDownloadArgs{
Server: source.Server,
Protocol: source.Protocol,
Certificate: source.Certificate,
Alias: source.Alias,
Type: info.Type,
AutoUpdate: true,
Public: info.Public,
StoragePool: poolName,
ProjectName: projectName,
Budget: -1,
})
if err != nil {
logger.Error("Failed to update the image", logger.Ctx{"err": err, "fingerprint": fingerprint})
continue
}
hash = newInfo.Fingerprint
if hash == fingerprint {
logger.Debug("Image already up to date", logger.Ctx{"fingerprint": fingerprint})
continue
}
var newID int
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
newID, _, err = tx.GetImage(ctx, hash, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
logger.Error("Error loading image", logger.Ctx{"err": err, "fingerprint": hash})
continue
}
if info.Cached {
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.SetImageCachedAndLastUseDate(ctx, projectName, hash, info.LastUsedAt)
})
if err != nil {
logger.Error("Error setting cached flag and last use date", logger.Ctx{"err": err, "fingerprint": hash})
continue
}
} else {
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
err = tx.UpdateImageLastUseDate(ctx, projectName, hash, info.LastUsedAt)
if err != nil {
logger.Error("Error setting last use date", logger.Ctx{"err": err, "fingerprint": hash})
return err
}
return nil
})
if err != nil {
continue
}
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.MoveImageAlias(ctx, id, newID)
})
if err != nil {
logger.Error("Error moving aliases", logger.Ctx{"err": err, "fingerprint": hash})
continue
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.CopyDefaultImageProfiles(ctx, id, newID)
})
if err != nil {
logger.Error("Copying default profiles", logger.Ctx{"err": err, "fingerprint": hash})
}
if poolName != "" {
pool, err := storagePools.LoadByName(s, poolName)
if err != nil {
logger.Error("Error loading storage pool to delete image", logger.Ctx{"err": err, "pool": poolName, "fingerprint": fingerprint})
continue
}
err = pool.DeleteImage(fingerprint, op)
if err != nil {
logger.Error("Error deleting image from storage pool", logger.Ctx{"err": err, "pool": pool.Name(), "fingerprint": fingerprint})
continue
}
}
err = s.Authorizer.AddImage(s.ShutdownCtx, projectName, info.Fingerprint)
if err != nil {
logger.Error("Failed to add image to authorizer", logger.Ctx{"fingerprint": info.Fingerprint, "project": projectName, "error": err})
}
var requestor *api.EventLifecycleRequestor
if op != nil {
requestor = op.Requestor()
}
s.Events.SendLifecycle(projectName, lifecycle.ImageCreated.Event(info.Fingerprint, projectName, requestor, logger.Ctx{"type": info.Type}))
}
if hash == fingerprint {
setRefreshResult(false)
return nil, nil
}
fname := filepath.Join(s.OS.VarDir, "images", fingerprint)
if util.PathExists(fname) {
err = os.Remove(fname)
if err != nil {
logger.Error("Error deleting image file", logger.Ctx{"fingerprint": fingerprint, "file": fname, "err": err})
}
}
fname = filepath.Join(s.OS.VarDir, "images", fingerprint) + ".rootfs"
if util.PathExists(fname) {
err = os.Remove(fname)
if err != nil {
logger.Error("Error deleting image rootfs file", logger.Ctx{"fingerprint": fingerprint, "file": fname, "err": err})
}
}
setRefreshResult(true)
return newInfo, nil
}
func pruneExpiredImagesTask(d *Daemon) (task.Func, task.Schedule) {
f := func(ctx context.Context) {
s := d.State()
opRun := func(op *operations.Operation) error {
return pruneExpiredImages(ctx, s, op)
}
op, err := operations.OperationCreate(s, "", operations.OperationClassTask, operationtype.ImagesExpire, nil, nil, opRun, nil, nil, nil)
if err != nil {
logger.Error("Failed creating expired image prune operation", logger.Ctx{"err": err})
return
}
logger.Debug("Acquiring image task lock")
imageTaskMu.Lock()
defer imageTaskMu.Unlock()
logger.Debug("Acquired image task lock")
logger.Info("Pruning expired images")
err = op.Start()
if err != nil {
logger.Error("Failed starting expired image prune operation", logger.Ctx{"err": err})
return
}
err = op.Wait(ctx)
if err != nil {
logger.Error("Failed expiring images", logger.Ctx{"err": err})
return
}
logger.Info("Done pruning expired images")
}
f(context.Background())
first := true
schedule := func() (time.Duration, error) {
interval := 24 * time.Hour
if first {
first = false
return interval, task.ErrSkip
}
return interval, nil
}
return f, schedule
}
func pruneLeftoverImages(s *state.State) {
opRun := func(op *operations.Operation) error {
var storageImages string
err := s.DB.Node.Transaction(s.ShutdownCtx, func(ctx context.Context, tx *db.NodeTx) error {
nodeConfig, err := node.ConfigLoad(ctx, tx)
if err != nil {
return err
}
storageImages = nodeConfig.StorageImagesVolume()
return nil
})
if err != nil {
return err
}
if storageImages != "" {
poolName, _, err := daemonStorageSplitVolume(storageImages)
if err != nil {
return err
}
pool, err := storagePools.LoadByName(s, poolName)
if err != nil {
return err
}
if pool.Driver().Info().VolumeMultiNode {
return nil
}
}
var images []string
err = s.DB.Cluster.Transaction(s.ShutdownCtx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
images, err = tx.GetLocalImagesFingerprints(ctx)
return err
})
if err != nil {
return fmt.Errorf("Unable to retrieve the list of images: %w", err)
}
entries, err := os.ReadDir(internalUtil.VarPath("images"))
if err != nil {
return fmt.Errorf("Unable to list the images directory: %w", err)
}
for _, entry := range entries {
fp := strings.Split(entry.Name(), ".")[0]
if !slices.Contains(images, fp) {
err = os.RemoveAll(internalUtil.VarPath("images", entry.Name()))
if err != nil {
return fmt.Errorf("Unable to remove leftover image: %v: %w", entry.Name(), err)
}
logger.Debugf("Removed leftover image file: %s", entry.Name())
}
}
return nil
}
op, err := operations.OperationCreate(s, "", operations.OperationClassTask, operationtype.ImagesPruneLeftover, nil, nil, opRun, nil, nil, nil)
if err != nil {
logger.Error("Failed creating leftover image clean up operation", logger.Ctx{"err": err})
return
}
logger.Debug("Acquiring image task lock")
imageTaskMu.Lock()
defer imageTaskMu.Unlock()
logger.Debug("Acquired image task lock")
logger.Info("Cleaning up leftover image files")
err = op.Start()
if err != nil {
logger.Error("Failed starting leftover image clean up operation", logger.Ctx{"err": err})
return
}
err = op.Wait(s.ShutdownCtx)
if err != nil {
logger.Error("Failed cleaning up leftover image files", logger.Ctx{"err": err})
return
}
logger.Infof("Done cleaning up leftover image files")
}
func pruneExpiredImages(ctx context.Context, s *state.State, op *operations.Operation) error {
var err error
var projectsImageRemoteCacheExpiryDays map[string]int64
var allImages map[string][]dbCluster.Image
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
globalImageRemoteCacheExpiryDays := s.GlobalConfig.ImagesRemoteCacheExpiryDays()
dbProjects, err := dbCluster.GetProjects(ctx, tx.Tx())
if err != nil {
return err
}
projectsImageRemoteCacheExpiryDays = make(map[string]int64, len(dbProjects))
for _, p := range dbProjects {
p, err := p.ToAPI(ctx, tx.Tx())
if err != nil {
return err
}
if p.Config["images.remote_cache_expiry"] != "" {
expiry, err := strconv.ParseInt(p.Config["images.remote_cache_expiry"], 10, 64)
if err != nil {
return fmt.Errorf("Unable to fetch project configuration: %w", err)
}
projectsImageRemoteCacheExpiryDays[p.Name] = expiry
} else {
projectsImageRemoteCacheExpiryDays[p.Name] = globalImageRemoteCacheExpiryDays
}
}
cached := true
images, err := dbCluster.GetImages(ctx, tx.Tx(), dbCluster.ImageFilter{Cached: &cached})
if err != nil {
return fmt.Errorf("Failed getting images: %w", err)
}
allImages = make(map[string][]dbCluster.Image, len(images))
for _, image := range images {
allImages[image.Fingerprint] = append(allImages[image.Fingerprint], image)
}
return nil
})
if err != nil {
return fmt.Errorf("Unable to retrieve project names: %w", err)
}
for fingerprint, dbImages := range allImages {
select {
case <-ctx.Done():
return nil
default:
}
dbImagesDeleted := 0
for _, dbImage := range dbImages {
expiryDays := projectsImageRemoteCacheExpiryDays[dbImage.Project]
if expiryDays <= 0 {
continue
}
timestamp := dbImage.UploadDate
if !dbImage.LastUseDate.Time.IsZero() {
timestamp = dbImage.LastUseDate.Time
}
imageExpiry := timestamp.Add(time.Duration(expiryDays) * time.Hour * 24)
if imageExpiry.After(time.Now()) {
continue
}
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
return tx.DeleteImage(ctx, dbImage.ID)
})
if err != nil {
return fmt.Errorf("Error deleting image %q in project %q from database: %w", fingerprint, dbImage.Project, err)
}
dbImagesDeleted++
logger.Info("Deleted expired cached image record", logger.Ctx{"fingerprint": fingerprint, "project": dbImage.Project, "expiry": imageExpiry})
s.Events.SendLifecycle(dbImage.Project, lifecycle.ImageDeleted.Event(fingerprint, dbImage.Project, op.Requestor(), nil))
}
if dbImagesDeleted < len(dbImages) {
continue
}
var poolIDs []int64
var poolNames []string
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
poolIDs, err = tx.GetPoolsWithImage(ctx, fingerprint)
if err != nil {
return err
}
poolNames, err = tx.GetPoolNamesFromIDs(ctx, poolIDs)
if err != nil {
return err
}
return nil
})
if err != nil {
continue
}
for _, poolName := range poolNames {
pool, err := storagePools.LoadByName(s, poolName)
if err != nil {
return fmt.Errorf("Error loading storage pool %q to delete image volume %q: %w", poolName, fingerprint, err)
}
err = pool.DeleteImage(fingerprint, op)
if err != nil {
return fmt.Errorf("Error deleting image volume %q from storage pool %q: %w", fingerprint, pool.Name(), err)
}
}
fname := filepath.Join(s.OS.VarDir, "images", fingerprint)
err = os.Remove(fname)
if err != nil && !os.IsNotExist(err) {
return fmt.Errorf("Error deleting image file %q: %w", fname, err)
}
fname = filepath.Join(s.OS.VarDir, "images", fingerprint) + ".rootfs"
err = os.Remove(fname)
if err != nil && !os.IsNotExist(err) {
return fmt.Errorf("Error deleting image file %q: %w", fname, err)
}
logger.Info("Deleted expired cached image files and volumes", logger.Ctx{"fingerprint": fingerprint})
}
return nil
}
func imageDelete(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var imgID int
var imgInfo *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
imgID, imgInfo, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
do := func(op *operations.Operation) error {
unlock, err := imageOperationLock(context.TODO(), imgInfo.Fingerprint)
if err != nil {
return err
}
defer unlock()
var exist bool
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
exist, err = tx.ImageExists(ctx, projectName, imgInfo.Fingerprint)
return err
})
if err != nil {
return err
}
if !exist {
return api.StatusErrorf(http.StatusNotFound, "Image not found")
}
if !isClusterNotification(r) {
var referenced bool
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
referenced, err = tx.ImageIsReferencedByOtherProjects(ctx, projectName, imgInfo.Fingerprint)
if err != nil {
return err
}
if referenced {
err = tx.DeleteImage(ctx, imgID)
if err != nil {
return fmt.Errorf("Error deleting image info from the database: %w", err)
}
}
return nil
})
if err != nil {
return err
}
if referenced {
return nil
}
notifier, err := cluster.NewNotifier(s, s.Endpoints.NetworkCert(), s.ServerCert(), cluster.NotifyAll)
if err != nil {
return err
}
err = notifier(func(client incus.InstanceServer) error {
op, err := client.UseProject(projectName).DeleteImage(imgInfo.Fingerprint)
if err != nil {
return fmt.Errorf("Failed to request to delete image from peer node: %w", err)
}
err = op.Wait()
if err != nil {
return fmt.Errorf("Failed to delete image from peer node: %w", err)
}
return nil
})
if err != nil {
return err
}
for _, alias := range imgInfo.Aliases {
err = s.Authorizer.DeleteImageAlias(s.ShutdownCtx, projectName, alias.Name)
if err != nil {
logger.Error("Failed to remove image alias from authorizer", logger.Ctx{"name": alias.Name, "project": projectName, "error": err})
}
s.Events.SendLifecycle(projectName, lifecycle.ImageAliasDeleted.Event(alias.Name, projectName, op.Requestor(), nil))
}
err = s.Authorizer.DeleteImage(s.ShutdownCtx, projectName, imgInfo.Fingerprint)
if err != nil {
logger.Error("Failed to remove image from authorizer", logger.Ctx{"fingerprint": imgInfo.Fingerprint, "project": projectName, "error": err})
}
s.Events.SendLifecycle(projectName, lifecycle.ImageDeleted.Event(imgInfo.Fingerprint, projectName, op.Requestor(), nil))
}
var poolIDs []int64
var poolNames []string
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
poolIDs, err = tx.GetPoolsWithImage(ctx, imgInfo.Fingerprint)
if err != nil {
return err
}
poolNames, err = tx.GetPoolNamesFromIDs(ctx, poolIDs)
if err != nil {
return err
}
return nil
})
if err != nil {
return err
}
for _, poolName := range poolNames {
pool, err := storagePools.LoadByName(s, poolName)
if err != nil {
return fmt.Errorf("Error loading storage pool %q to delete image %q: %w", poolName, imgInfo.Fingerprint, err)
}
if !isClusterNotification(r) || !pool.Driver().Info().Remote {
err = pool.DeleteImage(imgInfo.Fingerprint, op)
if err != nil {
return fmt.Errorf("Error deleting image %q from storage pool %q: %w", imgInfo.Fingerprint, pool.Name(), err)
}
}
}
if !isClusterNotification(r) {
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
return tx.DeleteImage(ctx, imgID)
})
if err != nil {
return fmt.Errorf("Error deleting image info from the database: %w", err)
}
}
imageDeleteFromDisk(imgInfo.Fingerprint)
return nil
}
resources := map[string][]api.URL{}
resources["images"] = []api.URL{*api.NewURL().Path(version.APIVersion, "images", imgInfo.Fingerprint)}
op, err := operations.OperationCreate(s, projectName, operations.OperationClassTask, operationtype.ImageDelete, resources, nil, do, nil, nil, r)
if err != nil {
return response.InternalError(err)
}
return operations.OperationResponse(op)
}
func imageDeleteFromDisk(fingerprint string) {
fname := internalUtil.VarPath("images", fingerprint)
if util.PathExists(fname) {
err := os.Remove(fname)
if err != nil && !os.IsNotExist(err) {
logger.Errorf("Error deleting image file %s: %s", fname, err)
}
}
fname = internalUtil.VarPath("images", fingerprint) + ".rootfs"
if util.PathExists(fname) {
err := os.Remove(fname)
if err != nil && !os.IsNotExist(err) {
logger.Errorf("Error deleting image file %s: %s", fname, err)
}
}
}
func doImageGet(ctx context.Context, tx *db.ClusterTx, project, fingerprint string, public bool) (*api.Image, error) {
filter := dbCluster.ImageFilter{Project: &project}
if public {
filter.Public = &public
}
_, imgInfo, err := tx.GetImageByFingerprintPrefix(ctx, fingerprint, filter)
if err != nil {
return nil, err
}
return imgInfo, nil
}
func imageValidSecret(s *state.State, r *http.Request, projectName string, fingerprint string, secret string) (*api.Operation, error) {
ops, err := operationsGetByType(s, r, projectName, operationtype.ImageToken)
if err != nil {
return nil, fmt.Errorf("Failed getting image token operations: %w", err)
}
for _, op := range ops {
if op.Resources == nil {
continue
}
opImages, ok := op.Resources["images"]
if !ok {
continue
}
if !util.StringPrefixInSlice(api.NewURL().Path(version.APIVersion, "images", fingerprint).String(), opImages) {
continue
}
opSecret, ok := op.Metadata["secret"]
if !ok {
continue
}
if opSecret == secret {
err = operationCancel(s, r, projectName, op)
if err != nil {
return nil, fmt.Errorf("Failed to cancel operation %q: %w", op.ID, err)
}
return op, nil
}
}
return nil, nil
}
func imageGet(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var info *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
info, err = doImageGet(ctx, tx, projectName, fingerprint, false)
if err != nil {
return err
}
return nil
})
if err != nil {
return response.SmartError(err)
}
var userCanViewImage bool
err = s.Authorizer.CheckPermission(r.Context(), r, auth.ObjectImage(projectName, info.Fingerprint), auth.EntitlementCanView)
if err == nil {
userCanViewImage = true
} else if !api.StatusErrorCheck(err, http.StatusForbidden) {
return response.SmartError(err)
}
public := d.checkTrustedClient(r) != nil || !userCanViewImage
secret := r.FormValue("secret")
op, err := imageValidSecret(s, r, projectName, info.Fingerprint, secret)
if err != nil {
return response.SmartError(err)
}
if !info.Public && public && op == nil {
return response.NotFound(fmt.Errorf("Image %q not found", info.Fingerprint))
}
etag := []any{info.Public, info.AutoUpdate, info.Properties}
return response.SyncResponseETag(true, info, etag)
}
func imagePut(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var id int
var info *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
id, info, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
etag := []any{info.Public, info.AutoUpdate, info.Properties}
err = localUtil.EtagCheck(r, etag)
if err != nil {
return response.PreconditionFailed(err)
}
req := api.ImagePut{}
err = json.NewDecoder(r.Body).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
if !req.ExpiresAt.IsZero() {
info.ExpiresAt = req.ExpiresAt
}
if req.Profiles == nil {
req.Profiles = []string{"default"}
}
profileIds := make([]int64, len(req.Profiles))
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
for i, profile := range req.Profiles {
profileID, _, err := tx.GetProfile(ctx, projectName, profile)
if response.IsNotFoundError(err) {
return fmt.Errorf("Profile '%s' doesn't exist", profile)
} else if err != nil {
return err
}
profileIds[i] = profileID
}
return tx.UpdateImage(ctx, id, info.Filename, info.Size, req.Public, req.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, req.Properties, projectName, profileIds)
})
if err != nil {
if response.IsNotFoundError(err) {
return response.BadRequest(err)
}
return response.SmartError(err)
}
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageUpdated.Event(info.Fingerprint, projectName, requestor, nil))
return response.EmptySyncResponse
}
func imagePatch(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var id int
var info *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
id, info, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
etag := []any{info.Public, info.AutoUpdate, info.Properties}
err = localUtil.EtagCheck(r, etag)
if err != nil {
return response.PreconditionFailed(err)
}
body, err := io.ReadAll(r.Body)
if err != nil {
return response.InternalError(err)
}
rdr1 := io.NopCloser(bytes.NewBuffer(body))
rdr2 := io.NopCloser(bytes.NewBuffer(body))
reqRaw := jmap.Map{}
err = json.NewDecoder(rdr1).Decode(&reqRaw)
if err != nil {
return response.BadRequest(err)
}
req := api.ImagePut{}
err = json.NewDecoder(rdr2).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
autoUpdate, err := reqRaw.GetBool("auto_update")
if err == nil {
info.AutoUpdate = autoUpdate
}
public, err := reqRaw.GetBool("public")
if err == nil {
info.Public = public
}
_, ok := reqRaw["properties"]
if ok {
properties := req.Properties
for k, v := range info.Properties {
_, ok := req.Properties[k]
if !ok {
properties[k] = v
}
}
info.Properties = properties
}
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
return tx.UpdateImage(ctx, id, info.Filename, info.Size, info.Public, info.AutoUpdate, info.Architecture, info.CreatedAt, info.ExpiresAt, info.Properties, "", nil)
})
if err != nil {
return response.SmartError(err)
}
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageUpdated.Event(info.Fingerprint, projectName, requestor, nil))
return response.EmptySyncResponse
}
func imageAliasesPost(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
req := api.ImageAliasesPost{}
err := json.NewDecoder(r.Body).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
if req.Name == "" || req.Target == "" {
return response.BadRequest(fmt.Errorf("name and target are required"))
}
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, _, err = tx.GetImageAlias(ctx, projectName, req.Name, true)
if !response.IsNotFoundError(err) {
if err != nil {
return err
}
return api.StatusErrorf(http.StatusConflict, "Alias %q already exists", req.Name)
}
imgID, _, err := tx.GetImageByFingerprintPrefix(ctx, req.Target, dbCluster.ImageFilter{Project: &projectName})
if err != nil {
return err
}
err = tx.CreateImageAlias(ctx, projectName, req.Name, imgID, req.Description)
if err != nil {
return err
}
return err
})
if err != nil {
return response.SmartError(err)
}
err = s.Authorizer.AddImageAlias(r.Context(), projectName, req.Name)
if err != nil {
logger.Error("Failed to add image alias to authorizer", logger.Ctx{"name": req.Name, "project": projectName, "error": err})
}
requestor := request.CreateRequestor(r)
lc := lifecycle.ImageAliasCreated.Event(req.Name, projectName, requestor, logger.Ctx{"target": req.Target})
s.Events.SendLifecycle(projectName, lc)
return response.SyncResponseLocation(true, nil, lc.Source)
}
func imageAliasesGet(d *Daemon, r *http.Request) response.Response {
projectName := request.ProjectParam(r)
recursion := localUtil.IsRecursionRequest(r)
s := d.State()
userHasPermission, err := s.Authorizer.GetPermissionChecker(r.Context(), r, auth.EntitlementCanView, auth.ObjectTypeImageAlias)
if err != nil {
return response.InternalError(fmt.Errorf("Failed to get a permission checker: %w", err))
}
var responseStr []string
var responseMap []api.ImageAliasesEntry
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
names, err := tx.GetImageAliases(ctx, projectName)
if err != nil {
return err
}
if recursion {
responseMap = make([]api.ImageAliasesEntry, 0, len(names))
} else {
responseStr = make([]string, 0, len(names))
}
for _, name := range names {
if !userHasPermission(auth.ObjectImageAlias(projectName, name)) {
continue
}
if !recursion {
responseStr = append(responseStr, api.NewURL().Path(version.APIVersion, "images", "aliases", name).String())
} else {
_, alias, err := tx.GetImageAlias(ctx, projectName, name, true)
if err != nil {
continue
}
responseMap = append(responseMap, alias)
}
}
return nil
})
if err != nil {
return response.SmartError(err)
}
if !recursion {
return response.SyncResponse(true, responseStr)
}
return response.SyncResponse(true, responseMap)
}
func imageAliasGet(d *Daemon, r *http.Request) response.Response {
projectName := request.ProjectParam(r)
name, err := url.PathUnescape(mux.Vars(r)["name"])
if err != nil {
return response.SmartError(err)
}
s := d.State()
var userCanViewImageAlias bool
err = s.Authorizer.CheckPermission(r.Context(), r, auth.ObjectImageAlias(projectName, name), auth.EntitlementCanView)
if err == nil {
userCanViewImageAlias = true
} else if !api.StatusErrorCheck(err, http.StatusForbidden) {
return response.SmartError(err)
}
public := d.checkTrustedClient(r) != nil || !userCanViewImageAlias
var alias api.ImageAliasesEntry
err = d.State().DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, alias, err = tx.GetImageAlias(ctx, projectName, name, !public)
return err
})
if err != nil {
return response.SmartError(err)
}
return response.SyncResponseETag(true, alias, alias)
}
func imageAliasDelete(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
name, err := url.PathUnescape(mux.Vars(r)["name"])
if err != nil {
return response.SmartError(err)
}
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, _, err = tx.GetImageAlias(ctx, projectName, name, true)
if err != nil {
return err
}
err = tx.DeleteImageAlias(ctx, projectName, name)
if err != nil {
return err
}
return nil
})
if err != nil {
return response.SmartError(err)
}
err = s.Authorizer.DeleteImageAlias(r.Context(), projectName, name)
if err != nil {
logger.Error("Failed to remove image alias from authorizer", logger.Ctx{"name": name, "project": projectName, "error": err})
}
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageAliasDeleted.Event(name, projectName, requestor, nil))
return response.EmptySyncResponse
}
func imageAliasPut(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
name, err := url.PathUnescape(mux.Vars(r)["name"])
if err != nil {
return response.SmartError(err)
}
req := api.ImageAliasesEntryPut{}
err = json.NewDecoder(r.Body).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
if req.Target == "" {
return response.BadRequest(fmt.Errorf("The target field is required"))
}
var imgAlias api.ImageAliasesEntry
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
var imgAliasID int
imgAliasID, imgAlias, err = tx.GetImageAlias(ctx, projectName, name, true)
if err != nil {
return err
}
err = localUtil.EtagCheck(r, imgAlias)
if err != nil {
return err
}
imageID, _, err := tx.GetImageByFingerprintPrefix(ctx, req.Target, dbCluster.ImageFilter{Project: &projectName})
if err != nil {
return err
}
err = tx.UpdateImageAlias(ctx, imgAliasID, imageID, req.Description)
if err != nil {
return err
}
return err
})
if err != nil {
return response.SmartError(err)
}
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageAliasUpdated.Event(imgAlias.Name, projectName, requestor, logger.Ctx{"target": req.Target}))
return response.EmptySyncResponse
}
func imageAliasPatch(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
name, err := url.PathUnescape(mux.Vars(r)["name"])
if err != nil {
return response.SmartError(err)
}
req := jmap.Map{}
err = json.NewDecoder(r.Body).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
var imgAlias api.ImageAliasesEntry
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
var imgAliasID int
imgAliasID, imgAlias, err = tx.GetImageAlias(ctx, projectName, name, true)
if err != nil {
return err
}
err = localUtil.EtagCheck(r, imgAlias)
if err != nil {
return err
}
_, ok := req["target"]
if ok {
target, err := req.GetString("target")
if err != nil {
return api.StatusErrorf(http.StatusBadRequest, "%v", err)
}
imgAlias.Target = target
}
_, ok = req["description"]
if ok {
description, err := req.GetString("description")
if err != nil {
return api.StatusErrorf(http.StatusBadRequest, "%v", err)
}
imgAlias.Description = description
}
imageID, _, err := tx.GetImage(ctx, imgAlias.Target, dbCluster.ImageFilter{Project: &projectName})
if err != nil {
return err
}
err = tx.UpdateImageAlias(ctx, imgAliasID, imageID, imgAlias.Description)
if err != nil {
return err
}
return nil
})
if err != nil {
return response.SmartError(err)
}
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageAliasUpdated.Event(imgAlias.Name, projectName, requestor, logger.Ctx{"target": imgAlias.Target}))
return response.EmptySyncResponse
}
func imageAliasPost(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
name, err := url.PathUnescape(mux.Vars(r)["name"])
if err != nil {
return response.SmartError(err)
}
req := api.ImageAliasesEntryPost{}
err = json.NewDecoder(r.Body).Decode(&req)
if err != nil {
return response.BadRequest(err)
}
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, _, err := tx.GetImageAlias(ctx, projectName, req.Name, true)
if !response.IsNotFoundError(err) {
if err != nil {
return err
}
return api.StatusErrorf(http.StatusConflict, "Alias %q already exists", req.Name)
}
imgAliasID, _, err := tx.GetImageAlias(ctx, projectName, name, true)
if err != nil {
return err
}
return tx.RenameImageAlias(ctx, imgAliasID, req.Name)
})
if err != nil {
return response.SmartError(err)
}
err = s.Authorizer.RenameImageAlias(r.Context(), projectName, name, req.Name)
if err != nil {
logger.Error("Failed to rename image alias in authorizer", logger.Ctx{"old_name": name, "new_name": req.Name, "project": projectName})
}
requestor := request.CreateRequestor(r)
lc := lifecycle.ImageAliasRenamed.Event(req.Name, projectName, requestor, logger.Ctx{"old_name": name})
s.Events.SendLifecycle(projectName, lc)
return response.SyncResponseLocation(true, nil, lc.Source)
}
func imageExport(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var imgInfo *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, imgInfo, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
var userCanViewImage bool
err = s.Authorizer.CheckPermission(r.Context(), r, auth.ObjectImage(projectName, imgInfo.Fingerprint), auth.EntitlementCanView)
if err == nil {
userCanViewImage = true
} else if !api.StatusErrorCheck(err, http.StatusForbidden) {
return response.SmartError(err)
}
public := d.checkTrustedClient(r) != nil || !userCanViewImage
secret := r.FormValue("secret")
if r.RemoteAddr == "@dev_incus" {
if imgInfo.Fingerprint != fingerprint {
return response.NotFound(fmt.Errorf("Image %q not found", fingerprint))
}
if !imgInfo.Public && !imgInfo.Cached {
return response.NotFound(fmt.Errorf("Image %q not found", fingerprint))
}
} else {
op, err := imageValidSecret(s, r, projectName, imgInfo.Fingerprint, secret)
if err != nil {
return response.SmartError(err)
}
if !imgInfo.Public && public && op == nil {
return response.NotFound(fmt.Errorf("Image %q not found", imgInfo.Fingerprint))
}
}
var address string
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
address, err = tx.LocateImage(ctx, imgInfo.Fingerprint)
return err
})
if err != nil {
return response.SmartError(err)
}
if address != "" {
client, err := cluster.Connect(address, s.Endpoints.NetworkCert(), s.ServerCert(), r, false)
if err != nil {
return response.SmartError(err)
}
return response.ForwardedResponse(client, r)
}
headers := map[string]string{}
headers["X-Incus-Type"] = "incus"
if imgInfo.Properties != nil && imgInfo.Properties["type"] == "oci" {
headers["X-Incus-Type"] = "oci"
}
imagePath := internalUtil.VarPath("images", imgInfo.Fingerprint)
rootfsPath := imagePath + ".rootfs"
_, ext, _, err := archive.DetectCompression(imagePath)
if err != nil {
ext = ""
}
filename := fmt.Sprintf("%s%s", imgInfo.Fingerprint, ext)
if util.PathExists(rootfsPath) {
files := make([]response.FileResponseEntry, 2)
files[0].Identifier = "metadata"
files[0].Path = imagePath
files[0].Filename = "meta-" + filename
_, ext, _, err = archive.DetectCompression(rootfsPath)
if err != nil {
ext = ""
}
filename = fmt.Sprintf("%s%s", imgInfo.Fingerprint, ext)
if imgInfo.Type == "virtual-machine" {
files[1].Identifier = "rootfs.img"
} else {
files[1].Identifier = "rootfs"
}
files[1].Path = rootfsPath
files[1].Filename = filename
return response.FileResponse(r, files, headers)
}
files := make([]response.FileResponseEntry, 1)
files[0].Identifier = filename
files[0].Path = imagePath
files[0].Filename = filename
requestor := request.CreateRequestor(r)
s.Events.SendLifecycle(projectName, lifecycle.ImageRetrieved.Event(imgInfo.Fingerprint, projectName, requestor, nil))
return response.FileResponse(r, files, headers)
}
func imageExportPost(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, _, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
req := api.ImageExportPost{}
err = json.NewDecoder(r.Body).Decode(&req)
if err != nil {
response.SmartError(err)
}
args := &incus.ConnectionArgs{
TLSServerCert: req.Certificate,
UserAgent: version.UserAgent,
Proxy: s.Proxy,
CachePath: s.OS.CacheDir,
CacheExpiry: time.Hour,
}
remote, err := incus.ConnectIncus(req.Target, args)
if err != nil {
return response.SmartError(err)
}
var imageCreateOp incus.Operation
run := func(op *operations.Operation) error {
createArgs := &incus.ImageCreateArgs{}
imageMetaPath := internalUtil.VarPath("images", fingerprint)
imageRootfsPath := internalUtil.VarPath("images", fingerprint+".rootfs")
metaFile, err := os.Open(imageMetaPath)
if err != nil {
return err
}
defer func() { _ = metaFile.Close() }()
createArgs.MetaFile = metaFile
createArgs.MetaName = filepath.Base(imageMetaPath)
if util.PathExists(imageRootfsPath) {
rootfsFile, err := os.Open(imageRootfsPath)
if err != nil {
return err
}
defer func() { _ = rootfsFile.Close() }()
createArgs.RootfsFile = rootfsFile
createArgs.RootfsName = filepath.Base(imageRootfsPath)
}
image := api.ImagesPost{
Filename: createArgs.MetaName,
Source: &api.ImagesPostSource{
Fingerprint: fingerprint,
Secret: req.Secret,
Mode: "push",
},
ImagePut: api.ImagePut{
Profiles: req.Profiles,
},
}
if req.Project != "" {
remote = remote.UseProject(req.Project)
}
imageCreateOp, err = remote.CreateImage(image, createArgs)
if err != nil {
return err
}
opAPI := imageCreateOp.Get()
var secret string
val, ok := opAPI.Metadata["secret"]
if ok {
secret = val.(string)
}
opWaitAPI, _, err := remote.GetOperationWaitSecret(opAPI.ID, secret, -1)
if err != nil {
return err
}
if opWaitAPI.StatusCode != api.Success {
return fmt.Errorf("Failed operation %q: %q", opWaitAPI.Status, opWaitAPI.Err)
}
s.Events.SendLifecycle(projectName, lifecycle.ImageRetrieved.Event(fingerprint, projectName, op.Requestor(), logger.Ctx{"target": req.Target}))
return nil
}
op, err := operations.OperationCreate(s, projectName, operations.OperationClassTask, operationtype.ImageDownload, nil, nil, run, nil, nil, r)
if err != nil {
return response.InternalError(err)
}
return operations.OperationResponse(op)
}
func imageSecret(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var imgInfo *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
_, imgInfo, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
return createTokenResponse(s, r, projectName, imgInfo.Fingerprint, nil)
}
func imageImportFromNode(imagesDir string, client incus.InstanceServer, fingerprint string) error {
buildDir, err := os.MkdirTemp(imagesDir, "incus_build_")
if err != nil {
return fmt.Errorf("failed to create temporary directory for download: %w", err)
}
defer func() { _ = os.RemoveAll(buildDir) }()
metaFile, err := os.CreateTemp(buildDir, "incus_tar_")
if err != nil {
return err
}
defer func() { _ = metaFile.Close() }()
rootfsFile, err := os.CreateTemp(buildDir, "incus_tar_")
if err != nil {
return err
}
defer func() { _ = rootfsFile.Close() }()
getReq := incus.ImageFileRequest{
MetaFile: io.WriteSeeker(metaFile),
RootfsFile: io.WriteSeeker(rootfsFile),
}
getResp, err := client.GetImageFile(fingerprint, getReq)
if err != nil {
return err
}
if getResp.RootfsSize > 0 {
err = rootfsFile.Truncate(getResp.RootfsSize)
if err != nil {
return err
}
}
err = metaFile.Truncate(getResp.MetaSize)
if err != nil {
return err
}
if getResp.RootfsSize == 0 {
rootfsPath := filepath.Join(imagesDir, fingerprint)
err := internalUtil.FileMove(metaFile.Name(), rootfsPath)
if err != nil {
return err
}
} else {
metaPath := filepath.Join(imagesDir, fingerprint)
rootfsPath := filepath.Join(imagesDir, fingerprint+".rootfs")
err := internalUtil.FileMove(metaFile.Name(), metaPath)
if err != nil {
return nil
}
err = internalUtil.FileMove(rootfsFile.Name(), rootfsPath)
if err != nil {
return nil
}
}
return nil
}
func imageRefresh(d *Daemon, r *http.Request) response.Response {
s := d.State()
projectName := request.ProjectParam(r)
fingerprint, err := url.PathUnescape(mux.Vars(r)["fingerprint"])
if err != nil {
return response.SmartError(err)
}
var imageID int
var imageInfo *api.Image
err = s.DB.Cluster.Transaction(r.Context(), func(ctx context.Context, tx *db.ClusterTx) error {
imageID, imageInfo, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &projectName})
return err
})
if err != nil {
return response.SmartError(err)
}
run := func(op *operations.Operation) error {
var nodes []string
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
nodes, err = tx.GetNodesWithImageAndAutoUpdate(ctx, fingerprint, true)
return err
})
if err != nil {
return fmt.Errorf("Error getting cluster members for refreshing image %q in project %q: %w", fingerprint, projectName, err)
}
newImage, err := autoUpdateImage(context.TODO(), s, op, imageID, imageInfo, projectName, true)
if err != nil {
return fmt.Errorf("Failed to update image %q in project %q: %w", fingerprint, projectName, err)
}
if newImage != nil {
if len(nodes) > 1 {
err := distributeImage(context.TODO(), s, nodes, fingerprint, newImage)
if err != nil {
return fmt.Errorf("Failed to distribute new image %q: %w", newImage.Fingerprint, err)
}
}
err = s.DB.Cluster.Transaction(context.TODO(), func(ctx context.Context, tx *db.ClusterTx) error {
return tx.DeleteImage(ctx, imageID)
})
if err != nil {
logger.Error("Error deleting old image from database", logger.Ctx{"err": err, "fingerprint": fingerprint, "ID": imageID})
}
}
return err
}
op, err := operations.OperationCreate(s, projectName, operations.OperationClassTask, operationtype.ImageRefresh, nil, nil, run, nil, nil, r)
if err != nil {
return response.InternalError(err)
}
return operations.OperationResponse(op)
}
func autoSyncImagesTask(s *state.State) (task.Func, task.Schedule) {
f := func(ctx context.Context) {
localClusterAddress := s.LocalConfig.ClusterAddress()
leader, err := s.Cluster.LeaderAddress()
if err != nil {
if errors.Is(err, cluster.ErrNodeIsNotClustered) {
return
}
logger.Error("Failed to get leader cluster member address", logger.Ctx{"err": err})
return
}
if localClusterAddress != leader {
logger.Debug("Skipping image synchronization task since we're not leader")
return
}
opRun := func(op *operations.Operation) error {
return autoSyncImages(ctx, s)
}
op, err := operations.OperationCreate(s, "", operations.OperationClassTask, operationtype.ImagesSynchronize, nil, nil, opRun, nil, nil, nil)
if err != nil {
logger.Error("Failed creating image synchronization operation", logger.Ctx{"err": err})
return
}
logger.Debug("Acquiring image task lock")
imageTaskMu.Lock()
defer imageTaskMu.Unlock()
logger.Debug("Acquired image task lock")
logger.Info("Synchronizing images across the cluster")
err = op.Start()
if err != nil {
logger.Error("Failed starting image synchronization operation", logger.Ctx{"err": err})
return
}
err = op.Wait(ctx)
if err != nil {
logger.Error("Failed synchronizing images", logger.Ctx{"err": err})
return
}
logger.Info("Done synchronizing images across the cluster")
}
return f, task.Hourly()
}
func autoSyncImages(ctx context.Context, s *state.State) error {
var imageProjectInfo map[string][]string
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
var err error
imageProjectInfo, err = tx.GetImages(ctx)
return err
})
if err != nil {
return fmt.Errorf("Failed to query image fingerprints: %w", err)
}
for fingerprint, projects := range imageProjectInfo {
ch := make(chan error)
go func(projectName string, fingerprint string) {
err := imageSyncBetweenNodes(ctx, s, nil, projectName, fingerprint)
if err != nil {
logger.Error("Failed to synchronize images", logger.Ctx{"err": err, "project": projectName, "fingerprint": fingerprint})
}
ch <- nil
}(projects[0], fingerprint)
select {
case <-ctx.Done():
return nil
case <-ch:
}
}
return nil
}
func imageSyncBetweenNodes(ctx context.Context, s *state.State, r *http.Request, project string, fingerprint string) error {
logger.Info("Syncing image to members started", logger.Ctx{"fingerprint": fingerprint, "project": project})
defer logger.Info("Syncing image to members finished", logger.Ctx{"fingerprint": fingerprint, "project": project})
var desiredSyncNodeCount int64
var syncNodeAddresses []string
err := s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
desiredSyncNodeCount = s.GlobalConfig.ImagesMinimalReplica()
if desiredSyncNodeCount == -1 {
nodesCount, err := tx.GetNodesCount(ctx)
if err != nil {
return fmt.Errorf("Failed to get the number of nodes: %w", err)
}
desiredSyncNodeCount = int64(nodesCount)
}
var err error
syncNodeAddresses, err = tx.GetNodesWithImage(ctx, fingerprint)
if err != nil {
return fmt.Errorf("Failed to get nodes for the image synchronization: %w", err)
}
return nil
})
if err != nil {
return err
}
if len(syncNodeAddresses) == 0 {
logger.Info("No members have image, nothing to do", logger.Ctx{"fingerprint": fingerprint, "project": project})
return nil
}
nodeCount := desiredSyncNodeCount - int64(len(syncNodeAddresses))
if nodeCount <= 0 {
logger.Info("Sufficient members have image", logger.Ctx{"fingerprint": fingerprint, "project": project, "desiredSyncCount": desiredSyncNodeCount, "syncedCount": len(syncNodeAddresses)})
return nil
}
syncNodeAddress := syncNodeAddresses[rand.Intn(len(syncNodeAddresses))]
source, err := cluster.Connect(syncNodeAddress, s.Endpoints.NetworkCert(), s.ServerCert(), r, true)
if err != nil {
return fmt.Errorf("Failed to connect to source node for image synchronization: %w", err)
}
source = source.UseProject(project)
var image *api.Image
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
_, image, err = tx.GetImage(ctx, fingerprint, dbCluster.ImageFilter{Project: &project})
return err
})
if err != nil {
return fmt.Errorf("Failed to get image: %w", err)
}
req := api.ImagesPost{
Source: &api.ImagesPostSource{
Fingerprint: image.Fingerprint,
Mode: "pull",
Type: "image",
Project: project,
},
}
for i := 0; i < int(nodeCount); i++ {
var addresses []string
err = s.DB.Cluster.Transaction(ctx, func(ctx context.Context, tx *db.ClusterTx) error {
addresses, err = tx.GetNodesWithoutImage(ctx, fingerprint)
return err
})
if err != nil {
return fmt.Errorf("Failed to get nodes for the image synchronization: %w", err)
}
if len(addresses) <= 0 {
logger.Info("All members have image", logger.Ctx{"fingerprint": fingerprint, "project": project})
return nil
}
targetNodeAddress := addresses[rand.Intn(len(addresses))]
client, err := cluster.Connect(targetNodeAddress, s.Endpoints.NetworkCert(), s.ServerCert(), r, true)
if err != nil {
return fmt.Errorf("Failed to connect node for image synchronization: %w", err)
}
client = client.UseProject(project)
logger.Info("Copying image to member", logger.Ctx{"fingerprint": fingerprint, "address": targetNodeAddress, "project": project})
op, err := client.CreateImage(req, nil)
if err != nil {
return fmt.Errorf("Failed to copy image to %q: %w", targetNodeAddress, err)
}
err = op.Wait()
if err != nil {
return err
}
}
return nil
}
func createTokenResponse(s *state.State, r *http.Request, projectName string, fingerprint string, metadata jmap.Map) response.Response {
secret, err := internalUtil.RandomHexString(32)
if err != nil {
return response.InternalError(err)
}
meta := jmap.Map{}
for k, v := range metadata {
meta[k] = v
}
meta["secret"] = secret
resources := map[string][]api.URL{}
resources["images"] = []api.URL{*api.NewURL().Path(version.APIVersion, "images", fingerprint)}
op, err := operations.OperationCreate(s, projectName, operations.OperationClassToken, operationtype.ImageToken, resources, meta, nil, nil, nil, r)
if err != nil {
return response.InternalError(err)
}
s.Events.SendLifecycle(projectName, lifecycle.ImageSecretCreated.Event(fingerprint, projectName, op.Requestor(), nil))
return operations.OperationResponse(op)
}