Merge branch 'feature-37'
This commit is contained in:
commit
d042a1ab5b
32 changed files with 1629 additions and 165 deletions
|
|
@ -42,27 +42,69 @@ jobs:
|
|||
goarch: ${{ matrix.goarch }}
|
||||
binari: ${{ matrix.binaries }}
|
||||
secrets: inherit
|
||||
upload-scripts:
|
||||
# Scripts et units systemd publiés comme assets de release : deploy.sh les
|
||||
# installe depuis la release, plus depuis la branche. Ajouter une entrée ici
|
||||
# suffit à livrer un nouveau fichier.
|
||||
upload-assets:
|
||||
runs-on: docker
|
||||
needs: [set-release-target]
|
||||
strategy:
|
||||
matrix:
|
||||
script:
|
||||
- run-dnsmasq-in-netns.sh
|
||||
include:
|
||||
- path: scripts/run-dnsmasq-in-netns.sh
|
||||
name: run-dnsmasq-in-netns.sh
|
||||
- path: systemd/agent.service
|
||||
name: agent.service
|
||||
- path: systemd/dnsmasq@.service
|
||||
name: dnsmasq@.service
|
||||
- path: systemd/metadata@.service
|
||||
name: metadata@.service
|
||||
steps:
|
||||
- uses: actions/checkout@v3
|
||||
- name: Move asset
|
||||
run: |
|
||||
mkdir -p "dist"
|
||||
cp scripts/${{ matrix.script }} dist/
|
||||
- name: Upload script
|
||||
cp "${{ matrix.path }}" dist/
|
||||
- name: Upload asset
|
||||
uses: actions/upload-artifact@v3
|
||||
with:
|
||||
name: ${{ matrix.script }}-${{ needs.set-release-target.outputs.release_cible }}
|
||||
path: dist/${{ matrix.script }}
|
||||
name: ${{ matrix.name }}-${{ needs.set-release-target.outputs.release_cible }}
|
||||
path: dist/${{ matrix.name }}
|
||||
# Manifeste des sommes de contrôle de tous les artefacts, au format sha256sum.
|
||||
# Dépend de tous les jobs qui produisent des artefacts : ajouter un producteur
|
||||
# sans l'ajouter ici donnerait un manifeste incomplet, donc un déploiement qui
|
||||
# refuse des assets légitimes.
|
||||
checksums:
|
||||
runs-on: docker
|
||||
needs: [set-release-target, build, upload-assets]
|
||||
steps:
|
||||
- name: Download all artifacts
|
||||
uses: actions/download-artifact@v3
|
||||
with:
|
||||
path: artifacts/
|
||||
- name: Générer SHA256SUMS
|
||||
run: |
|
||||
mkdir -p dist
|
||||
# download-artifact place chaque artefact dans son propre
|
||||
# sous-répertoire : on aplatit pour que le manifeste porte les noms
|
||||
# d'assets, tels que deploy.sh les demandera.
|
||||
find artifacts/ -type f -exec cp {} dist/ \;
|
||||
# Généré depuis dist/ pour que les noms soient nus, et hors de dist/
|
||||
# pour que le manifeste ne se liste pas lui-même.
|
||||
( cd dist && sha256sum * ) > SHA256SUMS
|
||||
cat SHA256SUMS
|
||||
- name: Upload SHA256SUMS
|
||||
uses: actions/upload-artifact@v3
|
||||
with:
|
||||
name: SHA256SUMS-${{ needs.set-release-target.outputs.release_cible }}
|
||||
path: SHA256SUMS
|
||||
prerelease:
|
||||
runs-on: docker
|
||||
needs: [set-release-target, build]
|
||||
# Tous les producteurs d'artefacts sont dans les needs : sans ça, le job
|
||||
# release peut appeler download-artifact avant la fin des uploads et publier
|
||||
# une release incomplète (assets ou manifeste manquants de façon non
|
||||
# déterministe).
|
||||
needs: [set-release-target, build, upload-assets, checksums]
|
||||
uses: ./.forgejo/workflows/release.yml
|
||||
with:
|
||||
tag: ${{ needs.set-release-target.outputs.release_cible }}
|
||||
|
|
|
|||
|
|
@ -91,7 +91,7 @@ paths:
|
|||
"404":
|
||||
$ref: "#/components/responses/NotFound"
|
||||
"409":
|
||||
description: VPC not in a deletable state
|
||||
description: VPC not deletable — only running or error states can be deleted, and all its subnets must be deleted first
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
|
|
@ -146,7 +146,7 @@ paths:
|
|||
schema:
|
||||
$ref: "#/components/schemas/Error"
|
||||
"422":
|
||||
description: Subnet not found or not in created state
|
||||
description: Subnet not found, or not in creating/running state
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
|
|
@ -185,6 +185,12 @@ paths:
|
|||
$ref: "#/components/schemas/VM"
|
||||
"404":
|
||||
$ref: "#/components/responses/NotFound"
|
||||
"409":
|
||||
description: VM not stoppable — only running or error states can be stopped
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/Error"
|
||||
"500":
|
||||
$ref: "#/components/responses/InternalError"
|
||||
|
||||
|
|
@ -274,6 +280,12 @@ paths:
|
|||
$ref: "#/components/schemas/Subnet"
|
||||
"404":
|
||||
$ref: "#/components/responses/NotFound"
|
||||
"409":
|
||||
description: Subnet not deletable — only running or error states can be deleted
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/Error"
|
||||
"500":
|
||||
$ref: "#/components/responses/InternalError"
|
||||
|
||||
|
|
@ -314,8 +326,8 @@ components:
|
|||
example: vp-00001
|
||||
state:
|
||||
type: string
|
||||
enum: [creating, created, deleting, deleted]
|
||||
example: created
|
||||
enum: [creating, running, error, deleting, deleted]
|
||||
example: running
|
||||
cidr:
|
||||
type: string
|
||||
example: "10.0.0.0/16"
|
||||
|
|
@ -373,8 +385,8 @@ components:
|
|||
example: sn-00001
|
||||
state:
|
||||
type: string
|
||||
enum: [creating, created, deleting, deleted]
|
||||
example: created
|
||||
enum: [creating, running, error, deleting, deleted]
|
||||
example: running
|
||||
vpc:
|
||||
type: string
|
||||
example: vpc1
|
||||
|
|
@ -472,8 +484,8 @@ components:
|
|||
example: vm-00001
|
||||
state:
|
||||
type: string
|
||||
enum: [starting, started, stopping, stopped]
|
||||
example: started
|
||||
enum: [creating, running, error, deleting, deleted]
|
||||
example: running
|
||||
metadata_port:
|
||||
type: string
|
||||
example: "80"
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ import (
|
|||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
||||
"git.g3e.fr/syonad/two/internal/migration"
|
||||
agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/pkg/logger"
|
||||
|
|
@ -31,6 +32,13 @@ func main() {
|
|||
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
||||
defer db.Close()
|
||||
|
||||
// Avant tout démarrage de service : la DB peut porter l'ancien vocabulaire
|
||||
// d'états, et des ressources transitoires orphelines d'un arrêt précédent.
|
||||
if err := migration.MigrateStates(db, log.With(slog.String("component", "migration"))); err != nil {
|
||||
log.Error("failed to migrate states", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
q := worker.New(cfg.Worker.BufferSize)
|
||||
q.Start(cfg.Worker.Count)
|
||||
|
||||
|
|
|
|||
|
|
@ -68,7 +68,13 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string)
|
|||
func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) {
|
||||
cmd := dispatcher.DeleteSubnetCommand{Name: name}
|
||||
if err := s.dispatcher.Prepare(cmd); err != nil {
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
// 404 si la ressource n'existe pas, 409 si elle existe mais n'est pas
|
||||
// dans un état supprimable — même convention que /vpcs et /vms.
|
||||
if _, dbErr := kv.GetFromDB(s.db, "subnet/"+name+"/state"); dbErr != nil {
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
} else {
|
||||
w.WriteHeader(http.StatusConflict)
|
||||
}
|
||||
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ func TestListSubnets_Empty(t *testing.T) {
|
|||
|
||||
func TestListSubnets_WithData(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
kv.AddInDB(db, "subnet/sn-2/state", "creating")
|
||||
kv.AddInDB(db, "subnet/sn-2/vpc", "vpc-1")
|
||||
|
|
@ -55,7 +55,7 @@ func TestListSubnets_InvalidMethod(t *testing.T) {
|
|||
|
||||
func TestPostSubnet_Created(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := SubnetCreateRequest{
|
||||
Name: "sn-new",
|
||||
VPC: "vpc-1",
|
||||
|
|
@ -91,7 +91,7 @@ func TestPostSubnet_MissingFields(t *testing.T) {
|
|||
|
||||
func TestPostSubnet_IfaceTypeOptional(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := SubnetCreateRequest{
|
||||
Name: "sn-opt",
|
||||
VPC: "vpc-1",
|
||||
|
|
@ -126,8 +126,8 @@ func TestPostSubnet_VPCNotFound(t *testing.T) {
|
|||
|
||||
func TestPostSubnet_Duplicate(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-exist/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-exist/state", "running")
|
||||
req := SubnetCreateRequest{
|
||||
Name: "sn-exist",
|
||||
VPC: "vpc-1",
|
||||
|
|
@ -163,7 +163,7 @@ func TestPostSubnet_VPCDeleting(t *testing.T) {
|
|||
|
||||
func TestPostSubnet_BridgeMode_Success(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := SubnetCreateRequest{
|
||||
Name: "sn-br",
|
||||
VPC: "vpc-1",
|
||||
|
|
@ -190,7 +190,7 @@ func TestPostSubnet_BridgeMode_Success(t *testing.T) {
|
|||
|
||||
func TestPostSubnet_UnknownMode(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := SubnetCreateRequest{
|
||||
Name: "sn-1",
|
||||
VPC: "vpc-1",
|
||||
|
|
@ -219,7 +219,7 @@ func TestPostSubnet_InvalidBody(t *testing.T) {
|
|||
|
||||
func TestGetSubnet_Found(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
kv.AddInDB(db, "subnet/sn-1/cidr", "10.0.0.0/24")
|
||||
kv.AddInDB(db, "subnet/sn-1/interface_ip", "10.0.0.1")
|
||||
|
|
@ -231,7 +231,7 @@ func TestGetSubnet_Found(t *testing.T) {
|
|||
}
|
||||
var result Subnet
|
||||
json.NewDecoder(w.Body).Decode(&result)
|
||||
if result.Name != "sn-1" || result.State != "created" {
|
||||
if result.Name != "sn-1" || result.State != "running" {
|
||||
t.Errorf("résultat inattendu : %+v", result)
|
||||
}
|
||||
if result.VPC != "vpc-1" {
|
||||
|
|
@ -261,7 +261,7 @@ func TestGetSubnet_EmptyName(t *testing.T) {
|
|||
|
||||
func TestDeleteSubnet_Success(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-del/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-del/state", "running")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-del", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.SubnetByNameHandler(w, req)
|
||||
|
|
@ -275,6 +275,17 @@ func TestDeleteSubnet_Success(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestDeleteSubnet_ConflictWhileCreating(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-wip/state", "creating")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-wip", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.SubnetByNameHandler(w, req)
|
||||
if w.Code != http.StatusConflict {
|
||||
t.Errorf("attendu 409, obtenu %d: %s", w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteSubnet_NotFound(t *testing.T) {
|
||||
s, _ := newTestServer(t)
|
||||
req := httptest.NewRequest(http.MethodDelete, "/subnets/inexistant", nil)
|
||||
|
|
@ -287,7 +298,7 @@ func TestDeleteSubnet_NotFound(t *testing.T) {
|
|||
|
||||
func TestSubnetByName_InvalidMethod(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
req := httptest.NewRequest(http.MethodPut, "/subnets/sn-1", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.SubnetByNameHandler(w, req)
|
||||
|
|
|
|||
|
|
@ -15,7 +15,7 @@ import (
|
|||
|
||||
func TestVmFromDB_SingleDisk(t *testing.T) {
|
||||
entries := map[string]string{
|
||||
"vm/vm-1/state": "started",
|
||||
"vm/vm-1/state": "running",
|
||||
"vm/vm-1/subnet": "sn-1",
|
||||
"vm/vm-1/ip": "10.0.0.5",
|
||||
"vm/vm-1/metadata_port": "1234",
|
||||
|
|
@ -37,7 +37,7 @@ func TestVmFromDB_SingleDisk(t *testing.T) {
|
|||
|
||||
func TestVmFromDB_MultiDisk(t *testing.T) {
|
||||
entries := map[string]string{
|
||||
"vm/vm-2/state": "started",
|
||||
"vm/vm-2/state": "running",
|
||||
"vm/vm-2/subnet": "sn-1",
|
||||
"vm/vm-2/ip": "10.0.0.6",
|
||||
"vm/vm-2/metadata_port": "1235",
|
||||
|
|
@ -62,7 +62,7 @@ func TestVmFromDB_MultiDisk(t *testing.T) {
|
|||
func TestVmFromDB_SlotGap(t *testing.T) {
|
||||
// sdb absent — sda et sdc seulement
|
||||
entries := map[string]string{
|
||||
"vm/vm-3/state": "started",
|
||||
"vm/vm-3/state": "running",
|
||||
"vm/vm-3/subnet": "sn-1",
|
||||
"vm/vm-3/ip": "10.0.0.7",
|
||||
"vm/vm-3/metadata_port": "1236",
|
||||
|
|
@ -94,7 +94,7 @@ func TestVmFromDB_SlotGap(t *testing.T) {
|
|||
|
||||
func TestStartVM_MultiDisk(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
body, _ := json.Marshal(VMCreateRequest{
|
||||
|
|
@ -128,7 +128,7 @@ func TestStartVM_MultiDisk(t *testing.T) {
|
|||
|
||||
func TestStartVM_StorageReturnedInResponse(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
body, _ := json.Marshal(VMCreateRequest{
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ func TestListVpcs_Empty(t *testing.T) {
|
|||
|
||||
func TestListVpcs_WithData(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/v1/state", "created")
|
||||
kv.AddInDB(db, "vpc/v1/state", "running")
|
||||
kv.AddInDB(db, "vpc/v2/state", "creating")
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcsHandler(w, httptest.NewRequest(http.MethodGet, "/vpcs", nil))
|
||||
|
|
@ -104,7 +104,7 @@ func TestPostVpc_InvalidCIDR(t *testing.T) {
|
|||
|
||||
func TestPostVpc_Duplicate(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-exist/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-exist/state", "running")
|
||||
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-exist", CIDR: "10.0.0.0/16"})
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
|
||||
|
|
@ -126,7 +126,7 @@ func TestPostVpc_InvalidBody(t *testing.T) {
|
|||
|
||||
func TestGetVpc_Found(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := httptest.NewRequest(http.MethodGet, "/vpcs/vpc-1", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcByNameHandler(w, req)
|
||||
|
|
@ -135,7 +135,7 @@ func TestGetVpc_Found(t *testing.T) {
|
|||
}
|
||||
var result VPC
|
||||
json.NewDecoder(w.Body).Decode(&result)
|
||||
if result.Name != "vpc-1" || result.State != "created" {
|
||||
if result.Name != "vpc-1" || result.State != "running" {
|
||||
t.Errorf("résultat inattendu : %+v", result)
|
||||
}
|
||||
}
|
||||
|
|
@ -162,7 +162,7 @@ func TestGetVpc_EmptyName(t *testing.T) {
|
|||
|
||||
func TestDeleteVpc_Success(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-del/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-del/state", "running")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-del", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcByNameHandler(w, req)
|
||||
|
|
@ -186,10 +186,37 @@ func TestDeleteVpc_NotFound(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestDeleteVpc_ConflictWhileCreating(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-wip/state", "creating")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-wip", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcByNameHandler(w, req)
|
||||
if w.Code != http.StatusConflict {
|
||||
t.Errorf("attendu 409, obtenu %d: %s", w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteVpc_AllowedFromError(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-ko/state", "error")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-ko", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcByNameHandler(w, req)
|
||||
if w.Code != http.StatusAccepted {
|
||||
t.Fatalf("attendu 202, obtenu %d: %s", w.Code, w.Body.String())
|
||||
}
|
||||
var result VPC
|
||||
json.NewDecoder(w.Body).Decode(&result)
|
||||
if result.State != "deleting" {
|
||||
t.Errorf("state attendu deleting, obtenu %q", result.State)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-busy/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-busy/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy")
|
||||
req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-busy", nil)
|
||||
w := httptest.NewRecorder()
|
||||
|
|
@ -201,7 +228,7 @@ func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) {
|
|||
|
||||
func TestVpcByName_InvalidMethod(t *testing.T) {
|
||||
s, db := newTestServer(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
req := httptest.NewRequest(http.MethodPut, "/vpcs/vpc-1", nil)
|
||||
w := httptest.NewRecorder()
|
||||
s.VpcByNameHandler(w, req)
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ import (
|
|||
"time"
|
||||
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/worker"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
|
@ -13,6 +14,7 @@ import (
|
|||
type Command interface {
|
||||
Prepare(db *badger.DB, cfg *configuration.Config) error
|
||||
Execute(db *badger.DB, cfg *configuration.Config) error
|
||||
Key() string
|
||||
}
|
||||
|
||||
type Dispatcher struct {
|
||||
|
|
@ -43,6 +45,10 @@ func (d *Dispatcher) Dispatch(cmd Command) {
|
|||
}
|
||||
if err != nil {
|
||||
d.logger.Error("command failed", append(attrs, "error", err)...)
|
||||
if setErr := state.Set(d.db, cmd.Key(), state.Error); setErr != nil {
|
||||
d.logger.Error("failed to mark resource as errored",
|
||||
"command", cmdType, "key", cmd.Key(), "error", setErr)
|
||||
}
|
||||
} else {
|
||||
d.logger.Info("command done", attrs...)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,8 +4,10 @@ import (
|
|||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
|
|
@ -61,3 +63,66 @@ func TestDispatcher_Dispatch_ExecuteErrorLogged(t *testing.T) {
|
|||
d.Dispatch(cmd)
|
||||
wg.Wait() // Execute s'est terminé — l'erreur est loggée, pas propagée
|
||||
}
|
||||
|
||||
func TestDispatcher_Dispatch_ExecuteErrorSetsErrorState(t *testing.T) {
|
||||
d, db := newTestDispatcher(t)
|
||||
if err := state.Set(db, "vpc/vpc-1", state.Creating); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
cmd := mockCmd{
|
||||
key: "vpc/vpc-1",
|
||||
prepareFn: func(*badger.DB, *configuration.Config) error { return nil },
|
||||
executeFn: func(*badger.DB, *configuration.Config) error {
|
||||
return errors.New("execute failed")
|
||||
},
|
||||
}
|
||||
d.Dispatch(cmd)
|
||||
|
||||
// L'état est écrit après le retour d'Execute : on scrute la DB.
|
||||
go func() {
|
||||
defer close(done)
|
||||
for {
|
||||
if s, err := state.Get(db, "vpc/vpc-1"); err == nil && s == state.Error {
|
||||
return
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
}
|
||||
}()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
s, _ := state.Get(db, "vpc/vpc-1")
|
||||
t.Fatalf("la ressource devrait être en %q, obtenu %q", state.Error, s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatcher_Dispatch_ExecuteSuccessKeepsState(t *testing.T) {
|
||||
d, db := newTestDispatcher(t)
|
||||
if err := state.Set(db, "vpc/vpc-1", state.Running); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
cmd := mockCmd{
|
||||
key: "vpc/vpc-1",
|
||||
prepareFn: func(*badger.DB, *configuration.Config) error { return nil },
|
||||
executeFn: func(*badger.DB, *configuration.Config) error {
|
||||
defer wg.Done()
|
||||
return nil
|
||||
},
|
||||
}
|
||||
d.Dispatch(cmd)
|
||||
wg.Wait()
|
||||
time.Sleep(50 * time.Millisecond) // laisse le temps d'une écriture parasite
|
||||
|
||||
s, err := state.Get(db, "vpc/vpc-1")
|
||||
if err != nil {
|
||||
t.Fatalf("Get a échoué : %v", err)
|
||||
}
|
||||
if s != state.Running {
|
||||
t.Errorf("l'état devrait rester %q, obtenu %q", state.Running, s)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -25,10 +25,18 @@ func newTestDispatcher(t *testing.T) (*Dispatcher, *badger.DB) {
|
|||
|
||||
// mockCmd implémente Command sans aucune dépendance système.
|
||||
type mockCmd struct {
|
||||
key string
|
||||
prepareFn func(*badger.DB, *configuration.Config) error
|
||||
executeFn func(*badger.DB, *configuration.Config) error
|
||||
}
|
||||
|
||||
func (m mockCmd) Key() string {
|
||||
if m.key == "" {
|
||||
return "vpc/mock"
|
||||
}
|
||||
return m.key
|
||||
}
|
||||
|
||||
func (m mockCmd) Prepare(db *badger.DB, cfg *configuration.Config) error {
|
||||
return m.prepareFn(db, cfg)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ import (
|
|||
"time"
|
||||
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/internal/subnet"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
|
|
@ -22,6 +23,8 @@ type CreateSubnetCommand struct {
|
|||
DefaultRoute bool
|
||||
}
|
||||
|
||||
func (c CreateSubnetCommand) Key() string { return "subnet/" + c.Name }
|
||||
|
||||
func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error {
|
||||
if c.Mode == "" {
|
||||
c.Mode = "vxlan"
|
||||
|
|
@ -32,18 +35,18 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e
|
|||
if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil {
|
||||
return fmt.Errorf("subnet %q already exists", c.Name)
|
||||
}
|
||||
vpcState, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state")
|
||||
vpcState, err := state.Get(db, "vpc/"+c.VPC)
|
||||
if err != nil {
|
||||
return fmt.Errorf("vpc %q not found", c.VPC)
|
||||
}
|
||||
if vpcState == "deleting" || vpcState == "deleted" {
|
||||
if vpcState != state.Creating && vpcState != state.Running {
|
||||
return fmt.Errorf("vpc %q is %s", c.VPC, vpcState)
|
||||
}
|
||||
localIface, ok := cfg.Interfaces[c.IfaceType]
|
||||
if !ok {
|
||||
localIface = cfg.DefaultInterface
|
||||
}
|
||||
kv.AddInDB(db, "subnet/"+c.Name+"/state", "creating")
|
||||
state.Set(db, c.Key(), state.Creating)
|
||||
kv.AddInDB(db, "subnet/"+c.Name+"/vpc", c.VPC)
|
||||
kv.AddInDB(db, "subnet/"+c.Name+"/mode", c.Mode)
|
||||
kv.AddInDB(db, "subnet/"+c.Name+"/local_iface", localIface)
|
||||
|
|
@ -59,16 +62,19 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e
|
|||
func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
|
||||
for {
|
||||
state, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state")
|
||||
vpcState, err := state.Get(db, "vpc/"+c.VPC)
|
||||
if err != nil {
|
||||
return fmt.Errorf("vpc %q not found while waiting", c.VPC)
|
||||
}
|
||||
if state == "created" {
|
||||
if vpcState == state.Running {
|
||||
break
|
||||
}
|
||||
if vpcState != state.Creating {
|
||||
return fmt.Errorf("vpc %q is %s, cannot create subnet %q", c.VPC, vpcState, c.Name)
|
||||
}
|
||||
select {
|
||||
case <-timeout:
|
||||
return fmt.Errorf("timed out waiting for vpc %q to be created", c.VPC)
|
||||
return fmt.Errorf("timed out waiting for vpc %q to be running", c.VPC)
|
||||
case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
|
||||
}
|
||||
}
|
||||
|
|
@ -79,22 +85,28 @@ type DeleteSubnetCommand struct {
|
|||
Name string
|
||||
}
|
||||
|
||||
func (c DeleteSubnetCommand) Key() string { return "subnet/" + c.Name }
|
||||
|
||||
func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||
if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil {
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return fmt.Errorf("subnet %q not found", c.Name)
|
||||
}
|
||||
return kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting")
|
||||
if !state.CanDelete(current) {
|
||||
return fmt.Errorf("subnet %q cannot be deleted while %s", c.Name, current)
|
||||
}
|
||||
return state.Set(db, c.Key(), state.Deleting)
|
||||
}
|
||||
|
||||
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
||||
if err := subnet.DeleteSubnet(db, c.Name); err != nil {
|
||||
return err
|
||||
}
|
||||
state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state")
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state == "deleted" {
|
||||
if current == state.Deleted {
|
||||
kv.DeleteInDB(db, "subnet/"+c.Name)
|
||||
}
|
||||
return nil
|
||||
|
|
|
|||
|
|
@ -17,7 +17,7 @@ func testCfg() *configuration.Config {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_Success(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -37,7 +37,7 @@ func TestCreateSubnetCommand_Prepare_Success(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -51,7 +51,7 @@ func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "inconnu", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -65,8 +65,8 @@ func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing.
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_Duplicate(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-exist/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-exist/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-exist", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -113,7 +113,7 @@ func TestCreateSubnetCommand_Prepare_VPCDeleted(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -130,7 +130,7 @@ func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -150,7 +150,7 @@ func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -163,7 +163,7 @@ func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", Mode: "vlan",
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -175,7 +175,7 @@ func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -195,7 +195,7 @@ func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) {
|
|||
|
||||
func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "running")
|
||||
cmd := CreateSubnetCommand{
|
||||
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
|
||||
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
|
||||
|
|
@ -211,7 +211,7 @@ func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) {
|
|||
|
||||
func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-del/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-del/state", "running")
|
||||
cmd := DeleteSubnetCommand{Name: "sn-del"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare a échoué : %v", err)
|
||||
|
|
@ -222,6 +222,32 @@ func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestDeleteSubnetCommand_Prepare_RefusedWhileCreating(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-wip/state", "creating")
|
||||
cmd := DeleteSubnetCommand{Name: "sn-wip"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait refuser la suppression d'un subnet en creating")
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "subnet/sn-wip/state")
|
||||
if s != "creating" {
|
||||
t.Errorf("l'état ne devrait pas changer, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteSubnetCommand_Prepare_AllowedFromError(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-ko/state", "error")
|
||||
cmd := DeleteSubnetCommand{Name: "sn-ko"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare devrait accepter un subnet en error : %v", err)
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "subnet/sn-ko/state")
|
||||
if s != "deleting" {
|
||||
t.Errorf("state attendu deleting, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteSubnetCommand_Prepare_NotFound(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
cmd := DeleteSubnetCommand{Name: "sn-inexistant"}
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ import (
|
|||
"time"
|
||||
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/internal/vm"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
|
|
@ -30,22 +31,24 @@ type StartVMCommand struct {
|
|||
SSHKey string
|
||||
}
|
||||
|
||||
func (c StartVMCommand) Key() string { return "vm/" + c.Name }
|
||||
|
||||
func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||
if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil {
|
||||
return fmt.Errorf("vm %q already exists", c.Name)
|
||||
}
|
||||
subnetState, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state")
|
||||
subnetState, err := state.Get(db, "subnet/"+c.Subnet)
|
||||
if err != nil {
|
||||
return fmt.Errorf("subnet %q not found", c.Subnet)
|
||||
}
|
||||
if subnetState == "deleting" || subnetState == "deleted" {
|
||||
if subnetState != state.Creating && subnetState != state.Running {
|
||||
return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState)
|
||||
}
|
||||
port, err := allocateMetadataPort(db)
|
||||
if err != nil {
|
||||
return fmt.Errorf("allocate metadata port: %w", err)
|
||||
}
|
||||
kv.AddInDB(db, "vm/"+c.Name+"/state", "starting")
|
||||
state.Set(db, c.Key(), state.Creating)
|
||||
kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet)
|
||||
kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP)
|
||||
kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port))
|
||||
|
|
@ -91,16 +94,19 @@ func allocateMetadataPort(db *badger.DB) (int, error) {
|
|||
func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
|
||||
for {
|
||||
state, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state")
|
||||
subnetState, err := state.Get(db, "subnet/"+c.Subnet)
|
||||
if err != nil {
|
||||
return fmt.Errorf("subnet %q not found while waiting", c.Subnet)
|
||||
}
|
||||
if state == "created" {
|
||||
if subnetState == state.Running {
|
||||
break
|
||||
}
|
||||
if subnetState != state.Creating {
|
||||
return fmt.Errorf("subnet %q is %s, cannot start vm %q", c.Subnet, subnetState, c.Name)
|
||||
}
|
||||
select {
|
||||
case <-timeout:
|
||||
return fmt.Errorf("timed out waiting for subnet %q to be created", c.Subnet)
|
||||
return fmt.Errorf("timed out waiting for subnet %q to be running", c.Subnet)
|
||||
case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
|
||||
}
|
||||
}
|
||||
|
|
@ -111,22 +117,28 @@ type StopVMCommand struct {
|
|||
Name string
|
||||
}
|
||||
|
||||
func (c StopVMCommand) Key() string { return "vm/" + c.Name }
|
||||
|
||||
func (c StopVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||
if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err != nil {
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return fmt.Errorf("vm %q not found", c.Name)
|
||||
}
|
||||
return kv.AddInDB(db, "vm/"+c.Name+"/state", "stopping")
|
||||
if !state.CanDelete(current) {
|
||||
return fmt.Errorf("vm %q cannot be stopped while %s", c.Name, current)
|
||||
}
|
||||
return state.Set(db, c.Key(), state.Deleting)
|
||||
}
|
||||
|
||||
func (c StopVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||
if err := vm.StopVM(db, c.Name, cfg); err != nil {
|
||||
return err
|
||||
}
|
||||
state, err := kv.GetFromDB(db, "vm/"+c.Name+"/state")
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state == "stopped" {
|
||||
if current == state.Deleted {
|
||||
kv.DeleteInDB(db, "vm/"+c.Name)
|
||||
}
|
||||
return nil
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ import (
|
|||
|
||||
func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
cmd := StartVMCommand{
|
||||
|
|
@ -34,7 +34,7 @@ func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) {
|
|||
|
||||
func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
cmd := StartVMCommand{
|
||||
|
|
@ -66,7 +66,7 @@ func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) {
|
|||
|
||||
func TestStartVMCommand_Prepare_SlotGap(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
// sdb absent au boot — slot réservé pour hotplug
|
||||
|
|
@ -96,7 +96,7 @@ func TestStartVMCommand_Prepare_SlotGap(t *testing.T) {
|
|||
|
||||
func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
|
||||
|
||||
cmd := StartVMCommand{
|
||||
|
|
@ -114,9 +114,54 @@ func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// --- StopVMCommand.Prepare ---
|
||||
|
||||
func TestStopVMCommand_Prepare_Success(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vm/vm-run/state", "running")
|
||||
cmd := StopVMCommand{Name: "vm-run"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare a échoué : %v", err)
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "vm/vm-run/state")
|
||||
if s != "deleting" {
|
||||
t.Errorf("state attendu deleting, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStopVMCommand_Prepare_RefusedWhileCreating(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vm/vm-wip/state", "creating")
|
||||
cmd := StopVMCommand{Name: "vm-wip"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait refuser l'arrêt d'une VM en creating")
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "vm/vm-wip/state")
|
||||
if s != "creating" {
|
||||
t.Errorf("l'état ne devrait pas changer, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStopVMCommand_Prepare_AllowedFromError(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vm/vm-ko/state", "error")
|
||||
cmd := StopVMCommand{Name: "vm-ko"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare devrait accepter une VM en error : %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStopVMCommand_Prepare_NotFound(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
cmd := StopVMCommand{Name: "vm-inexistante"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait échouer si la VM n'existe pas")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStartVMCommand_Prepare_Duplicate(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vm/vm-exist/state", "started")
|
||||
kv.AddInDB(db, "vm/vm-exist/state", "running")
|
||||
|
||||
cmd := StartVMCommand{
|
||||
Name: "vm-exist",
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
"time"
|
||||
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/internal/vpc"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
|
|
@ -17,6 +18,8 @@ type CreateVPCCommand struct {
|
|||
CIDR string
|
||||
}
|
||||
|
||||
func (c CreateVPCCommand) Key() string { return "vpc/" + c.Name }
|
||||
|
||||
func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil {
|
||||
return fmt.Errorf("vpc %q already exists", c.Name)
|
||||
|
|
@ -27,7 +30,7 @@ func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error
|
|||
if err := kv.AddInDB(db, "vpc/"+c.Name+"/cidr", c.CIDR); err != nil {
|
||||
return err
|
||||
}
|
||||
return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
|
||||
return state.Set(db, c.Key(), state.Creating)
|
||||
}
|
||||
|
||||
func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
||||
|
|
@ -38,10 +41,16 @@ type DeleteVPCCommand struct {
|
|||
Name string
|
||||
}
|
||||
|
||||
func (c DeleteVPCCommand) Key() string { return "vpc/" + c.Name }
|
||||
|
||||
func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err != nil {
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return fmt.Errorf("vpc %q not found", c.Name)
|
||||
}
|
||||
if !state.CanDelete(current) {
|
||||
return fmt.Errorf("vpc %q cannot be deleted while %s", c.Name, current)
|
||||
}
|
||||
entries, err := kv.ListByPrefix(db, "subnet/")
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to list subnets: %w", err)
|
||||
|
|
@ -51,12 +60,12 @@ func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error
|
|||
continue
|
||||
}
|
||||
subnetName := strings.Split(key, "/")[1]
|
||||
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||
if err != nil || (state != "deleting" && state != "deleted") {
|
||||
s, err := state.Get(db, "subnet/"+subnetName)
|
||||
if err != nil || (s != state.Deleting && s != state.Deleted) {
|
||||
return fmt.Errorf("subnet %q must be deleted before deleting vpc %q", subnetName, c.Name)
|
||||
}
|
||||
}
|
||||
return kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting")
|
||||
return state.Set(db, c.Key(), state.Deleting)
|
||||
}
|
||||
|
||||
func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||
|
|
@ -70,8 +79,8 @@ func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) erro
|
|||
for key, value := range entries {
|
||||
if strings.HasSuffix(key, "/vpc") && value == c.Name {
|
||||
subnetName := strings.Split(key, "/")[1]
|
||||
state, _ := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||
if state == "deleting" {
|
||||
s, _ := state.Get(db, "subnet/"+subnetName)
|
||||
if s == state.Deleting {
|
||||
pending = true
|
||||
break
|
||||
}
|
||||
|
|
@ -89,11 +98,11 @@ func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) erro
|
|||
if err := vpc.DeleteVPC(db, c.Name); err != nil {
|
||||
return err
|
||||
}
|
||||
state, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state")
|
||||
current, err := state.Get(db, c.Key())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state == "deleted" {
|
||||
if current == state.Deleted {
|
||||
kv.DeleteInDB(db, "vpc/"+c.Name)
|
||||
}
|
||||
return nil
|
||||
|
|
|
|||
|
|
@ -32,7 +32,7 @@ func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
|
|||
|
||||
func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-exist/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-exist/state", "running")
|
||||
cmd := CreateVPCCommand{Name: "vpc-exist", CIDR: "10.0.0.0/16"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait échouer sur un VPC déjà existant")
|
||||
|
|
@ -51,7 +51,7 @@ func TestCreateVPCCommand_Prepare_InvalidCIDR(t *testing.T) {
|
|||
|
||||
func TestDeleteVPCCommand_Prepare_Success(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-del/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-del/state", "running")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-del"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare a échoué : %v", err)
|
||||
|
|
@ -70,10 +70,45 @@ func TestDeleteVPCCommand_Prepare_NotFound(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestDeleteVPCCommand_Prepare_RefusedWhileCreating(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-wip/state", "creating")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-wip"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait refuser la suppression d'un VPC en creating")
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "vpc/vpc-wip/state")
|
||||
if s != "creating" {
|
||||
t.Errorf("l'état ne devrait pas changer, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteVPCCommand_Prepare_RefusedWhileDeleting(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-gone/state", "deleting")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-gone"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
t.Error("Prepare devrait refuser un VPC déjà en deleting")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteVPCCommand_Prepare_AllowedFromError(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-ko/state", "error")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-ko"}
|
||||
if err := cmd.Prepare(db, nil); err != nil {
|
||||
t.Fatalf("Prepare devrait accepter un VPC en error : %v", err)
|
||||
}
|
||||
s, _ := kv.GetFromDB(db, "vpc/vpc-ko/state")
|
||||
if s != "deleting" {
|
||||
t.Errorf("state attendu deleting, obtenu %q", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeleteVPCCommand_Prepare_BlockedByActiveSubnet(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-busy/state", "created")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-busy/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-busy"}
|
||||
if err := cmd.Prepare(db, nil); err == nil {
|
||||
|
|
@ -83,7 +118,7 @@ func TestDeleteVPCCommand_Prepare_BlockedByActiveSubnet(t *testing.T) {
|
|||
|
||||
func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleted(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-ok/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-ok/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "deleted")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-ok"}
|
||||
|
|
@ -94,7 +129,7 @@ func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleted(t *testing.T) {
|
|||
|
||||
func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleting(t *testing.T) {
|
||||
_, db := newTestDispatcher(t)
|
||||
kv.AddInDB(db, "vpc/vpc-ok/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-ok/state", "running")
|
||||
kv.AddInDB(db, "subnet/sn-1/state", "deleting")
|
||||
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok")
|
||||
cmd := DeleteVPCCommand{Name: "vpc-ok"}
|
||||
|
|
|
|||
84
internal/migration/state.go
Normal file
84
internal/migration/state.go
Normal file
|
|
@ -0,0 +1,84 @@
|
|||
// Package migration met la DB au format attendu par la version courante de
|
||||
// l'agent. Les migrations sont idempotentes et jouées au démarrage.
|
||||
package migration
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
// prefixes énumère les familles de ressources portant une clé /state.
|
||||
var prefixes = []string{"vpc/", "subnet/", "vm/"}
|
||||
|
||||
// legacyStates traduit l'ancien vocabulaire (états VPC/subnet d'avant
|
||||
// l'unification, et états VM start/stop) vers l'enum courant.
|
||||
var legacyStates = map[string]state.State{
|
||||
"created": state.Running,
|
||||
"started": state.Running,
|
||||
"starting": state.Creating,
|
||||
"stopping": state.Deleting,
|
||||
"stopped": state.Deleted,
|
||||
}
|
||||
|
||||
// MigrateStates convertit les valeurs de state héritées vers l'enum courant,
|
||||
// puis marque en erreur les ressources restées dans un état transitoire.
|
||||
//
|
||||
// La réconciliation est sûre au démarrage : la queue worker est en mémoire et
|
||||
// vide à ce moment, donc aucune commande n'est en cours. Une ressource en
|
||||
// creating/deleting est nécessairement orpheline d'un arrêt de l'agent — sans
|
||||
// ce passage en error, elle resterait indéfiniment non supprimable.
|
||||
func MigrateStates(db *badger.DB, log *slog.Logger) error {
|
||||
for _, prefix := range prefixes {
|
||||
entries, err := kv.ListByPrefix(db, prefix)
|
||||
if err != nil {
|
||||
return fmt.Errorf("list %s: %w", prefix, err)
|
||||
}
|
||||
for key, value := range entries {
|
||||
if !strings.HasSuffix(key, "/state") {
|
||||
continue
|
||||
}
|
||||
resource := strings.TrimSuffix(key, "/state")
|
||||
target, reason := targetState(value)
|
||||
if target == "" {
|
||||
continue
|
||||
}
|
||||
if err := state.Set(db, resource, target); err != nil {
|
||||
return fmt.Errorf("migrate %s: %w", key, err)
|
||||
}
|
||||
log.Info("state migrated",
|
||||
"resource", resource, "from", value, "to", string(target), "reason", reason)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// targetState retourne l'état vers lequel migrer une valeur brute, ou "" si
|
||||
// elle doit rester inchangée.
|
||||
func targetState(value string) (state.State, string) {
|
||||
if legacy, ok := legacyStates[value]; ok {
|
||||
// Un état hérité transitoire est orphelin au même titre qu'un état
|
||||
// transitoire courant : on applique la réconciliation directement.
|
||||
if state.IsTransient(legacy) {
|
||||
return state.Error, "orphaned legacy state"
|
||||
}
|
||||
return legacy, "legacy state"
|
||||
}
|
||||
|
||||
current, err := state.Parse(value)
|
||||
if err != nil {
|
||||
// Valeur inconnue : ni l'ancien vocabulaire, ni le nouveau. On la
|
||||
// bascule en error plutôt que de la laisser bloquer l'agent — la
|
||||
// ressource reste visible et supprimable.
|
||||
return state.Error, "unknown state"
|
||||
}
|
||||
if state.IsTransient(current) {
|
||||
return state.Error, "orphaned transient state"
|
||||
}
|
||||
return "", ""
|
||||
}
|
||||
175
internal/migration/state_test.go
Normal file
175
internal/migration/state_test.go
Normal file
|
|
@ -0,0 +1,175 @@
|
|||
package migration
|
||||
|
||||
import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"testing"
|
||||
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func newTestDB(t *testing.T) *badger.DB {
|
||||
t.Helper()
|
||||
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
|
||||
t.Cleanup(func() { db.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
func discardLogger() *slog.Logger {
|
||||
return slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||
}
|
||||
|
||||
func TestMigrateStates_LegacyMapping(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
cases := map[string]struct {
|
||||
resource string
|
||||
from string
|
||||
want state.State
|
||||
}{
|
||||
"vpc created": {"vpc/vpc-1", "created", state.Running},
|
||||
"subnet created": {"subnet/sn-1", "created", state.Running},
|
||||
"vm started": {"vm/vm-1", "started", state.Running},
|
||||
"vm stopped": {"vm/vm-2", "stopped", state.Deleted},
|
||||
// Les états transitoires hérités sont orphelins après un redémarrage.
|
||||
"vm starting": {"vm/vm-3", "starting", state.Error},
|
||||
"vm stopping": {"vm/vm-4", "stopping", state.Error},
|
||||
"vpc creating": {"vpc/vpc-2", "creating", state.Error},
|
||||
"subnet deleting": {"subnet/sn-2", "deleting", state.Error},
|
||||
}
|
||||
for _, c := range cases {
|
||||
kv.AddInDB(db, c.resource+"/state", c.from)
|
||||
}
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("MigrateStates a échoué : %v", err)
|
||||
}
|
||||
|
||||
for name, c := range cases {
|
||||
got, err := state.Get(db, c.resource)
|
||||
if err != nil {
|
||||
t.Errorf("%s : Get a échoué : %v", name, err)
|
||||
continue
|
||||
}
|
||||
if got != c.want {
|
||||
t.Errorf("%s : %q → %q, attendu %q", name, c.from, got, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateStates_StableStatesUntouched(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
stable := map[string]state.State{
|
||||
"vpc/vpc-1": state.Running,
|
||||
"subnet/sn-1": state.Deleted,
|
||||
"vm/vm-1": state.Error,
|
||||
}
|
||||
for resource, s := range stable {
|
||||
if err := state.Set(db, resource, s); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("MigrateStates a échoué : %v", err)
|
||||
}
|
||||
|
||||
for resource, want := range stable {
|
||||
got, err := state.Get(db, resource)
|
||||
if err != nil {
|
||||
t.Fatalf("Get(%s) a échoué : %v", resource, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Errorf("%s : état modifié %q, attendu %q", resource, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateStates_Idempotent(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vm/vm-1/state", "starting")
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("premier passage : %v", err)
|
||||
}
|
||||
first := map[string]state.State{}
|
||||
for _, r := range []string{"vpc/vpc-1", "vm/vm-1"} {
|
||||
s, err := state.Get(db, r)
|
||||
if err != nil {
|
||||
t.Fatalf("Get(%s) : %v", r, err)
|
||||
}
|
||||
first[r] = s
|
||||
}
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("second passage : %v", err)
|
||||
}
|
||||
for r, want := range first {
|
||||
got, err := state.Get(db, r)
|
||||
if err != nil {
|
||||
t.Fatalf("Get(%s) : %v", r, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Errorf("%s : le second passage a modifié l'état (%q → %q)", r, want, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateStates_UnknownValueBecomesError(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
kv.AddInDB(db, "vpc/vpc-corrompu/state", "n'importe quoi")
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("MigrateStates a échoué : %v", err)
|
||||
}
|
||||
|
||||
got, err := state.Get(db, "vpc/vpc-corrompu")
|
||||
if err != nil {
|
||||
t.Fatalf("Get a échoué : %v", err)
|
||||
}
|
||||
if got != state.Error {
|
||||
t.Errorf("une valeur inconnue devrait devenir %q, obtenu %q", state.Error, got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateStates_IgnoresNonStateKeys(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
kv.AddInDB(db, "vpc/vpc-1/state", "created")
|
||||
kv.AddInDB(db, "vpc/vpc-1/cidr", "10.0.0.0/16")
|
||||
kv.AddInDB(db, "subnet/sn-1/local_iface", "br-vms")
|
||||
kv.AddInDB(db, "vm/vm-1/disk/sda", "/data/root.qcow2")
|
||||
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("MigrateStates a échoué : %v", err)
|
||||
}
|
||||
|
||||
untouched := map[string]string{
|
||||
"vpc/vpc-1/cidr": "10.0.0.0/16",
|
||||
"subnet/sn-1/local_iface": "br-vms",
|
||||
"vm/vm-1/disk/sda": "/data/root.qcow2",
|
||||
}
|
||||
for key, want := range untouched {
|
||||
got, err := kv.GetFromDB(db, key)
|
||||
if err != nil {
|
||||
t.Errorf("la clé %s devrait exister : %v", key, err)
|
||||
continue
|
||||
}
|
||||
if got != want {
|
||||
t.Errorf("%s = %q, attendu %q", key, got, want)
|
||||
}
|
||||
}
|
||||
// Aucune clé /state parasite ne doit apparaître sur ces ressources.
|
||||
if _, err := kv.GetFromDB(db, "vm/vm-1/state"); err == nil {
|
||||
t.Error("aucune clé state ne devrait être créée pour vm-1")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateStates_EmptyDB(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
if err := MigrateStates(db, discardLogger()); err != nil {
|
||||
t.Fatalf("MigrateStates devrait réussir sur une DB vide : %v", err)
|
||||
}
|
||||
}
|
||||
|
|
@ -3,12 +3,13 @@ package agentmetrics
|
|||
import (
|
||||
"strings"
|
||||
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
var allStates = []string{"creating", "created", "deleting", "deleted"}
|
||||
var allStates = state.All()
|
||||
|
||||
// AgentCollector implements prometheus.Collector and exposes agent metrics
|
||||
// by querying the BadgerDB on each scrape.
|
||||
|
|
@ -16,6 +17,7 @@ type AgentCollector struct {
|
|||
db *badger.DB
|
||||
vpcsTotal *prometheus.Desc
|
||||
subnetsTotal *prometheus.Desc
|
||||
vmsTotal *prometheus.Desc
|
||||
}
|
||||
|
||||
func NewAgentCollector(db *badger.DB) *AgentCollector {
|
||||
|
|
@ -31,23 +33,30 @@ func NewAgentCollector(db *badger.DB) *AgentCollector {
|
|||
"Number of subnets by state.",
|
||||
[]string{"state"}, nil,
|
||||
),
|
||||
vmsTotal: prometheus.NewDesc(
|
||||
"syonad_vms_total",
|
||||
"Number of VMs by state.",
|
||||
[]string{"state"}, nil,
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *AgentCollector) Describe(ch chan<- *prometheus.Desc) {
|
||||
ch <- c.vpcsTotal
|
||||
ch <- c.subnetsTotal
|
||||
ch <- c.vmsTotal
|
||||
}
|
||||
|
||||
func (c *AgentCollector) Collect(ch chan<- prometheus.Metric) {
|
||||
c.collectStates(ch, "vpc/", c.vpcsTotal)
|
||||
c.collectStates(ch, "subnet/", c.subnetsTotal)
|
||||
c.collectStates(ch, "vm/", c.vmsTotal)
|
||||
}
|
||||
|
||||
// collectStates counts resources under the given DB prefix by their state value
|
||||
// and emits one gauge per state label.
|
||||
func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix string, desc *prometheus.Desc) {
|
||||
counts := make(map[string]float64, len(allStates))
|
||||
counts := make(map[state.State]float64, len(allStates))
|
||||
for _, s := range allStates {
|
||||
counts[s] = 0
|
||||
}
|
||||
|
|
@ -55,13 +64,18 @@ func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix strin
|
|||
items, err := kv.ListByPrefix(c.db, prefix)
|
||||
if err == nil {
|
||||
for key, val := range items {
|
||||
if strings.HasSuffix(key, "/state") {
|
||||
counts[val]++
|
||||
if !strings.HasSuffix(key, "/state") {
|
||||
continue
|
||||
}
|
||||
// Une valeur hors enum n'est pas comptée : la migration au
|
||||
// démarrage les a toutes ramenées dans l'enum.
|
||||
if s, err := state.Parse(val); err == nil {
|
||||
counts[s]++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for _, state := range allStates {
|
||||
ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[state], state)
|
||||
for _, s := range allStates {
|
||||
ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[s], string(s))
|
||||
}
|
||||
}
|
||||
|
|
|
|||
151
internal/prometheus/agent/collector_test.go
Normal file
151
internal/prometheus/agent/collector_test.go
Normal file
|
|
@ -0,0 +1,151 @@
|
|||
package agentmetrics
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
)
|
||||
|
||||
func newTestDB(t *testing.T) *badger.DB {
|
||||
t.Helper()
|
||||
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
|
||||
t.Cleanup(func() { db.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
// collect exécute une collecte et retourne, par nom de métrique, la valeur de
|
||||
// chaque série indexée par son label state.
|
||||
func collect(t *testing.T, c *AgentCollector) map[string]map[string]float64 {
|
||||
t.Helper()
|
||||
ch := make(chan prometheus.Metric, 64)
|
||||
c.Collect(ch)
|
||||
close(ch)
|
||||
|
||||
out := map[string]map[string]float64{}
|
||||
for m := range ch {
|
||||
var pb dto.Metric
|
||||
if err := m.Write(&pb); err != nil {
|
||||
t.Fatalf("write metric: %v", err)
|
||||
}
|
||||
// Le fqName n'est pas exposé directement : il figure dans la
|
||||
// représentation textuelle du Desc.
|
||||
desc := m.Desc().String()
|
||||
var name string
|
||||
switch {
|
||||
case strings.Contains(desc, "syonad_vpcs_total"):
|
||||
name = "vpcs"
|
||||
case strings.Contains(desc, "syonad_subnets_total"):
|
||||
name = "subnets"
|
||||
case strings.Contains(desc, "syonad_vms_total"):
|
||||
name = "vms"
|
||||
default:
|
||||
t.Fatalf("métrique inattendue : %s", desc)
|
||||
}
|
||||
if out[name] == nil {
|
||||
out[name] = map[string]float64{}
|
||||
}
|
||||
for _, l := range pb.GetLabel() {
|
||||
if l.GetName() == "state" {
|
||||
out[name][l.GetValue()] = pb.GetGauge().GetValue()
|
||||
}
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func TestCollector_CountsByState(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
for resource, s := range map[string]state.State{
|
||||
"vpc/vpc-1": state.Running,
|
||||
"vpc/vpc-2": state.Running,
|
||||
"vpc/vpc-3": state.Error,
|
||||
"subnet/sn-1": state.Creating,
|
||||
"vm/vm-1": state.Running,
|
||||
"vm/vm-2": state.Error,
|
||||
} {
|
||||
if err := state.Set(db, resource, s); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
got := collect(t, NewAgentCollector(db))
|
||||
wantVPC := map[string]float64{"creating": 0, "running": 2, "error": 1, "deleting": 0, "deleted": 0}
|
||||
for s, want := range wantVPC {
|
||||
if got["vpcs"][s] != want {
|
||||
t.Errorf("vpcs{state=%q} = %v, attendu %v", s, got["vpcs"][s], want)
|
||||
}
|
||||
}
|
||||
if got["subnets"]["creating"] != 1 {
|
||||
t.Errorf("subnets{state=creating} = %v, attendu 1", got["subnets"]["creating"])
|
||||
}
|
||||
for s, want := range map[string]float64{"running": 1, "error": 1, "deleting": 0} {
|
||||
if got["vms"][s] != want {
|
||||
t.Errorf("vms{state=%q} = %v, attendu %v", s, got["vms"][s], want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Les clés vm/<name>/disk/<dev> ne portent pas de suffixe /state : elles ne
|
||||
// doivent pas être comptées comme des ressources.
|
||||
func TestCollector_IgnoresVMDiskKeys(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
if err := state.Set(db, "vm/vm-1", state.Running); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
kv.AddInDB(db, "vm/vm-1/disk/sda", "/data/sda.qcow2")
|
||||
kv.AddInDB(db, "vm/vm-1/disk/vda", "/data/vda.qcow2")
|
||||
|
||||
got := collect(t, NewAgentCollector(db))
|
||||
if got["vms"]["running"] != 1 {
|
||||
t.Errorf("vms{state=running} = %v, attendu 1", got["vms"]["running"])
|
||||
}
|
||||
}
|
||||
|
||||
// Une série par état doit être émise même à zéro : une métrique qui disparaît
|
||||
// côté Prometheus casse les alertes qui s'en servent.
|
||||
func TestCollector_EmitsEveryState(t *testing.T) {
|
||||
got := collect(t, NewAgentCollector(newTestDB(t)))
|
||||
for _, family := range []string{"vpcs", "subnets", "vms"} {
|
||||
if len(got[family]) != len(state.All()) {
|
||||
t.Errorf("%s : %d séries, attendu %d", family, len(got[family]), len(state.All()))
|
||||
}
|
||||
for _, s := range state.All() {
|
||||
if _, ok := got[family][string(s)]; !ok {
|
||||
t.Errorf("%s : série manquante pour l'état %q", family, s)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollector_IgnoresNonStateKeys(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
if err := state.Set(db, "vpc/vpc-1", state.Running); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
kv.AddInDB(db, "vpc/vpc-1/cidr", "10.0.0.0/16")
|
||||
|
||||
got := collect(t, NewAgentCollector(db))
|
||||
if got["vpcs"]["running"] != 1 {
|
||||
t.Errorf("vpcs{state=running} = %v, attendu 1", got["vpcs"]["running"])
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollector_Describe(t *testing.T) {
|
||||
ch := make(chan *prometheus.Desc, 10)
|
||||
NewAgentCollector(newTestDB(t)).Describe(ch)
|
||||
close(ch)
|
||||
|
||||
var n int
|
||||
for range ch {
|
||||
n++
|
||||
}
|
||||
if n != 3 {
|
||||
t.Errorf("%d descripteurs, attendu 3", n)
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,9 @@
|
|||
package qemu
|
||||
|
||||
func ScopeName(vmName string) string {
|
||||
return "two-vm-" + vmName + ".scope"
|
||||
}
|
||||
|
||||
type DiskConfig struct {
|
||||
Path string
|
||||
Dev string
|
||||
|
|
|
|||
|
|
@ -103,9 +103,16 @@ func Start(cfg Config) error {
|
|||
"-daemonize",
|
||||
)
|
||||
|
||||
cmd := exec.Command("qemu-system-x86_64", args...)
|
||||
if err := cmd.Run(); err != nil {
|
||||
return fmt.Errorf("qemu-system-x86_64: %w", err)
|
||||
scopeArgs := append([]string{
|
||||
"--scope",
|
||||
"--unit=" + ScopeName(cfg.Name),
|
||||
"--collect",
|
||||
"qemu-system-x86_64",
|
||||
}, args...)
|
||||
|
||||
cmd := exec.Command("systemd-run", scopeArgs...)
|
||||
if out, err := cmd.CombinedOutput(); err != nil {
|
||||
return fmt.Errorf("systemd-run qemu-system-x86_64: %w: %s", err, strings.TrimSpace(string(out)))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
68
internal/state/state.go
Normal file
68
internal/state/state.go
Normal file
|
|
@ -0,0 +1,68 @@
|
|||
// Package state définit les états possibles d'une ressource (VPC, subnet, VM)
|
||||
// et centralise leur lecture/écriture en DB.
|
||||
package state
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
type State string
|
||||
|
||||
const (
|
||||
Creating State = "creating"
|
||||
Running State = "running"
|
||||
Error State = "error"
|
||||
Deleting State = "deleting"
|
||||
Deleted State = "deleted"
|
||||
)
|
||||
|
||||
// All retourne tous les états valides, dans l'ordre du cycle de vie.
|
||||
func All() []State {
|
||||
return []State{Creating, Running, Error, Deleting, Deleted}
|
||||
}
|
||||
|
||||
// Parse convertit une valeur brute lue en DB en State.
|
||||
func Parse(s string) (State, error) {
|
||||
for _, valid := range All() {
|
||||
if State(s) == valid {
|
||||
return valid, nil
|
||||
}
|
||||
}
|
||||
return "", fmt.Errorf("unknown state %q", s)
|
||||
}
|
||||
|
||||
// CanDelete indique si une ressource dans cet état peut être supprimée.
|
||||
// Une ressource en cours de création ou déjà en cours de suppression ne l'est pas ;
|
||||
// une ressource en erreur l'est, pour permettre le nettoyage d'un échec partiel.
|
||||
func CanDelete(s State) bool {
|
||||
return s == Running || s == Error
|
||||
}
|
||||
|
||||
// IsTransient indique si l'état suppose une commande en cours d'exécution.
|
||||
// Au démarrage de l'agent, la queue worker est vide : une ressource dans un état
|
||||
// transitoire est donc orpheline.
|
||||
func IsTransient(s State) bool {
|
||||
return s == Creating || s == Deleting
|
||||
}
|
||||
|
||||
// Get lit l'état d'une ressource. prefix est la racine de la ressource,
|
||||
// sans le suffixe /state (ex: "vpc/vpc-1").
|
||||
func Get(db *badger.DB, prefix string) (State, error) {
|
||||
raw, err := kv.GetFromDB(db, prefix+"/state")
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return Parse(raw)
|
||||
}
|
||||
|
||||
// Set écrit l'état d'une ressource. prefix suit la même convention que Get.
|
||||
func Set(db *badger.DB, prefix string, s State) error {
|
||||
if _, err := Parse(string(s)); err != nil {
|
||||
return err
|
||||
}
|
||||
return kv.AddInDB(db, prefix+"/state", string(s))
|
||||
}
|
||||
168
internal/state/state_test.go
Normal file
168
internal/state/state_test.go
Normal file
|
|
@ -0,0 +1,168 @@
|
|||
package state
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
// newTestDB ouvre une base BadgerDB dans un répertoire temporaire.
|
||||
// La base est fermée automatiquement en fin de test.
|
||||
func newTestDB(t *testing.T) *badger.DB {
|
||||
t.Helper()
|
||||
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
|
||||
t.Cleanup(func() { db.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
// --- Parse ---
|
||||
|
||||
func TestParse_ValidStates(t *testing.T) {
|
||||
for _, s := range All() {
|
||||
got, err := Parse(string(s))
|
||||
if err != nil {
|
||||
t.Errorf("Parse(%q) a échoué : %v", s, err)
|
||||
}
|
||||
if got != s {
|
||||
t.Errorf("Parse(%q) = %q, attendu %q", s, got, s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestParse_UnknownState(t *testing.T) {
|
||||
// "created" et "started" sont les anciens états : ils ne doivent plus être
|
||||
// acceptés une fois la migration passée.
|
||||
for _, raw := range []string{"created", "started", "stopping", "", "RUNNING"} {
|
||||
if _, err := Parse(raw); err == nil {
|
||||
t.Errorf("Parse(%q) devrait échouer", raw)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- All ---
|
||||
|
||||
func TestAll_ContainsEveryState(t *testing.T) {
|
||||
want := map[State]bool{Creating: false, Running: false, Error: false, Deleting: false, Deleted: false}
|
||||
for _, s := range All() {
|
||||
if _, ok := want[s]; !ok {
|
||||
t.Errorf("All() retourne un état inattendu %q", s)
|
||||
}
|
||||
want[s] = true
|
||||
}
|
||||
for s, seen := range want {
|
||||
if !seen {
|
||||
t.Errorf("All() ne contient pas %q", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- CanDelete ---
|
||||
|
||||
func TestCanDelete(t *testing.T) {
|
||||
cases := map[State]bool{
|
||||
Running: true,
|
||||
Error: true,
|
||||
Creating: false,
|
||||
Deleting: false,
|
||||
Deleted: false,
|
||||
}
|
||||
for s, want := range cases {
|
||||
if got := CanDelete(s); got != want {
|
||||
t.Errorf("CanDelete(%q) = %v, attendu %v", s, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- IsTransient ---
|
||||
|
||||
func TestIsTransient(t *testing.T) {
|
||||
cases := map[State]bool{
|
||||
Creating: true,
|
||||
Deleting: true,
|
||||
Running: false,
|
||||
Error: false,
|
||||
Deleted: false,
|
||||
}
|
||||
for s, want := range cases {
|
||||
if got := IsTransient(s); got != want {
|
||||
t.Errorf("IsTransient(%q) = %v, attendu %v", s, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- Get / Set ---
|
||||
|
||||
func TestSetGet_RoundTrip(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
|
||||
if err := Set(db, "vpc/vpc-1", Creating); err != nil {
|
||||
t.Fatalf("Set a échoué : %v", err)
|
||||
}
|
||||
got, err := Get(db, "vpc/vpc-1")
|
||||
if err != nil {
|
||||
t.Fatalf("Get a échoué : %v", err)
|
||||
}
|
||||
if got != Creating {
|
||||
t.Errorf("Get = %q, attendu %q", got, Creating)
|
||||
}
|
||||
|
||||
if err := Set(db, "vpc/vpc-1", Running); err != nil {
|
||||
t.Fatalf("Set a échoué : %v", err)
|
||||
}
|
||||
got, err = Get(db, "vpc/vpc-1")
|
||||
if err != nil {
|
||||
t.Fatalf("Get a échoué : %v", err)
|
||||
}
|
||||
if got != Running {
|
||||
t.Errorf("Get après écrasement = %q, attendu %q", got, Running)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSet_WritesStateSuffix(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
|
||||
if err := Set(db, "vm/vm-1", Running); err != nil {
|
||||
t.Fatalf("Set a échoué : %v", err)
|
||||
}
|
||||
raw, err := kv.GetFromDB(db, "vm/vm-1/state")
|
||||
if err != nil {
|
||||
t.Fatalf("la clé vm/vm-1/state devrait exister : %v", err)
|
||||
}
|
||||
if raw != "running" {
|
||||
t.Errorf("valeur en DB = %q, attendu \"running\"", raw)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSet_InvalidState(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
|
||||
if err := Set(db, "vpc/vpc-1", State("created")); err == nil {
|
||||
t.Error("Set devrait refuser un état inconnu")
|
||||
}
|
||||
if _, err := kv.GetFromDB(db, "vpc/vpc-1/state"); err == nil {
|
||||
t.Error("aucune clé ne devrait être écrite pour un état invalide")
|
||||
}
|
||||
}
|
||||
|
||||
func TestGet_MissingKey(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
|
||||
if _, err := Get(db, "vpc/inconnu"); err == nil {
|
||||
t.Error("Get devrait échouer sur une ressource absente")
|
||||
}
|
||||
}
|
||||
|
||||
func TestGet_CorruptedValue(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
|
||||
// Valeur écrite hors du package (ancien état, corruption) : Get doit
|
||||
// remonter l'erreur plutôt que de retourner un State silencieusement vide.
|
||||
if err := kv.AddInDB(db, "subnet/sn-1/state", "created"); err != nil {
|
||||
t.Fatalf("préparation du test : %v", err)
|
||||
}
|
||||
if _, err := Get(db, "subnet/sn-1"); err == nil {
|
||||
t.Error("Get devrait échouer sur une valeur non reconnue")
|
||||
}
|
||||
}
|
||||
|
|
@ -7,18 +7,18 @@ import (
|
|||
"git.g3e.fr/syonad/two/internal/ebtables"
|
||||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/systemd"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func CreateSubnet(db *badger.DB, subnetName string) error {
|
||||
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||
current, err := state.Get(db, "subnet/"+subnetName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state != "creating" {
|
||||
if current != state.Creating {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -31,7 +31,7 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
|
|||
return err
|
||||
}
|
||||
|
||||
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created")
|
||||
return state.Set(db, "subnet/"+subnetName, state.Running)
|
||||
}
|
||||
|
||||
func createSubnet(db *badger.DB, subnetName string, d subnetData) error {
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
"git.g3e.fr/syonad/two/internal/ebtables"
|
||||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/pkg/systemd"
|
||||
|
||||
|
|
@ -14,11 +15,11 @@ import (
|
|||
)
|
||||
|
||||
func DeleteSubnet(db *badger.DB, subnetName string) error {
|
||||
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||
current, err := state.Get(db, "subnet/"+subnetName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state != "deleting" {
|
||||
if current != state.Deleting {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -44,7 +45,7 @@ func DeleteSubnet(db *badger.DB, subnetName string) error {
|
|||
return fmt.Errorf("unknown subnet mode %q", d.mode)
|
||||
}
|
||||
|
||||
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted")
|
||||
return state.Set(db, "subnet/"+subnetName, state.Deleted)
|
||||
}
|
||||
|
||||
func stopDHCP(db *badger.DB, subnetName string, d subnetData) error {
|
||||
|
|
|
|||
|
|
@ -12,17 +12,17 @@ import (
|
|||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/internal/qemu"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
|
||||
state, err := kv.GetFromDB(db, "vm/"+name+"/state")
|
||||
current, err := state.Get(db, "vm/"+name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state != "starting" {
|
||||
if current != state.Creating {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -84,7 +84,7 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
|
|||
return fmt.Errorf("start qemu: %w", err)
|
||||
}
|
||||
|
||||
return kv.AddInDB(db, "vm/"+name+"/state", "started")
|
||||
return state.Set(db, "vm/"+name, state.Running)
|
||||
}
|
||||
|
||||
func copyFile(src, dst string) error {
|
||||
|
|
|
|||
|
|
@ -12,17 +12,17 @@ import (
|
|||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/internal/qmp"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
|
||||
state, err := kv.GetFromDB(db, "vm/"+name+"/state")
|
||||
current, err := state.Get(db, "vm/"+name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if state != "stopping" {
|
||||
if current != state.Deleting {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -65,7 +65,7 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
|
|||
os.Remove(varsPath)
|
||||
}
|
||||
|
||||
return kv.AddInDB(db, "vm/"+name+"/state", "stopped")
|
||||
return state.Set(db, "vm/"+name, state.Deleted)
|
||||
}
|
||||
|
||||
func waitQMPDead(socketPath string, timeout, poll time.Duration) {
|
||||
|
|
|
|||
|
|
@ -5,15 +5,15 @@ import (
|
|||
|
||||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func CreateVPC(db *badger.DB, name string) error {
|
||||
if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil {
|
||||
if current, err := state.Get(db, "vpc/"+name); err != nil {
|
||||
return err
|
||||
} else if state == "creating" {
|
||||
} else if current == state.Creating {
|
||||
vpcID := strings.SplitN(name, "-", 2)[1]
|
||||
|
||||
if err := netns.Create(name); err != nil {
|
||||
|
|
@ -48,7 +48,7 @@ func CreateVPC(db *badger.DB, name string) error {
|
|||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
kv.AddInDB(db, "vpc/"+name+"/state", "created")
|
||||
return state.Set(db, "vpc/"+name, state.Running)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,15 +5,15 @@ import (
|
|||
|
||||
"git.g3e.fr/syonad/two/internal/netif"
|
||||
"git.g3e.fr/syonad/two/internal/netns"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/internal/state"
|
||||
|
||||
"github.com/dgraph-io/badger/v4"
|
||||
)
|
||||
|
||||
func DeleteVPC(db *badger.DB, name string) error {
|
||||
if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil {
|
||||
if current, err := state.Get(db, "vpc/"+name); err != nil {
|
||||
return err
|
||||
} else if state == "deleting" {
|
||||
} else if current == state.Deleting {
|
||||
vpcID := strings.SplitN(name, "-", 2)[1]
|
||||
|
||||
if err := netif.DeleteLink("vp-" + vpcID + "-e"); err != nil {
|
||||
|
|
@ -23,7 +23,7 @@ func DeleteVPC(db *badger.DB, name string) error {
|
|||
if err := netns.Delete(name); err != nil {
|
||||
return err
|
||||
}
|
||||
kv.AddInDB(db, "vpc/"+name+"/state", "deleted")
|
||||
return state.Set(db, "vpc/"+name, state.Deleted)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
|
|||
169
scripts/bootstrap_kvm.sh
Executable file
169
scripts/bootstrap_kvm.sh
Executable file
|
|
@ -0,0 +1,169 @@
|
|||
#!/bin/bash
|
||||
#
|
||||
# Préparation d'un host au profil kvm : paquets, kernel, réseau de base.
|
||||
#
|
||||
# Ce fichier n'est pas exécutable seul : il est *sourcé* par deploy.sh quand
|
||||
# --bootstrap est demandé, et réutilise ses helpers (run, info, warn, die,
|
||||
# exec_with_dry_run, FLAGS_TRUE). Conséquences :
|
||||
# - aucun effet de bord au chargement : uniquement des définitions,
|
||||
# - aucune variable globale de deploy.sh redéfinie ici (SCRIPT_PATH, TAG,
|
||||
# BIN_PATH, ASSETS…) : le sourcing les écraserait,
|
||||
# - pas de `set -e` ni de garde `BASH_SOURCE` en fin de fichier.
|
||||
#
|
||||
# Contrat : tout bootstrap_<profil>.sh expose bootstrap_host() avec cette
|
||||
# signature. deploy.sh n'a rien à savoir du contenu du profil.
|
||||
#
|
||||
# L'host est stateless (root sur tmpfs) : rien de ce qui suit ne survit à un
|
||||
# reboot, deploy.sh --bootstrap est rejoué à chaque démarrage.
|
||||
|
||||
KVM_PACKAGES="qemu-system-x86 ovmf dnsmasq ebtables iptables nfs-common jq curl"
|
||||
|
||||
kvm_packages () {
|
||||
local DRY_RUN="${1}"
|
||||
local WITH_PACKAGES="${2}"
|
||||
|
||||
[[ ${WITH_PACKAGES} -eq ${FLAGS_TRUE} ]] || { info "paquets : ignorés (--nopackages)"; return 0; }
|
||||
|
||||
info "paquets"
|
||||
run "${DRY_RUN}" "DEBIAN_FRONTEND=noninteractive apt-get update"
|
||||
run "${DRY_RUN}" "DEBIAN_FRONTEND=noninteractive apt-get install -y ${KVM_PACKAGES}"
|
||||
|
||||
# Le dnsmasq système prendrait le port 53 en concurrence des instances
|
||||
# dnsmasq@ lancées par l'agent dans les netns.
|
||||
run "${DRY_RUN}" "systemctl disable --now dnsmasq.service || true"
|
||||
run "${DRY_RUN}" "systemctl mask dnsmasq.service"
|
||||
}
|
||||
|
||||
kvm_kernel () {
|
||||
local DRY_RUN="${1}"
|
||||
|
||||
info "kernel"
|
||||
|
||||
# modprobe AVANT sysctl : les clés net.bridge.* n'existent pas tant que
|
||||
# br_netfilter n'est pas chargé, et sysctl --system échouerait.
|
||||
run "${DRY_RUN}" "modprobe br_netfilter"
|
||||
run "${DRY_RUN}" "echo br_netfilter > /etc/modules-load.d/two.conf"
|
||||
|
||||
# bridge-nf-call-iptables est requis par le DNAT metadata (169.254.169.254,
|
||||
# cf. internal/iptables) : sans lui, iptables ne voit pas le trafic bridgé
|
||||
# des VMs. Contrepartie : tout le trafic inter-VM traverse les tables nat.
|
||||
run "${DRY_RUN}" "printf 'net.ipv4.ip_forward = 1\nnet.bridge.bridge-nf-call-iptables = 1\n' > /etc/sysctl.d/90-two.conf"
|
||||
run "${DRY_RUN}" "sysctl --system >/dev/null"
|
||||
}
|
||||
|
||||
# Bridge vide, up, sans STP.
|
||||
kvm_ensure_bridge () {
|
||||
local DRY_RUN="${1}"
|
||||
local NAME="${2}"
|
||||
|
||||
if ! ip link show dev "${NAME}" >/dev/null 2>&1
|
||||
then
|
||||
run "${DRY_RUN}" "ip link add name '${NAME}' type bridge"
|
||||
fi
|
||||
run "${DRY_RUN}" "ip link set dev '${NAME}' type bridge stp_state 0"
|
||||
run "${DRY_RUN}" "ip link set up dev '${NAME}'"
|
||||
}
|
||||
|
||||
kvm_network () {
|
||||
local DRY_RUN="${1}"
|
||||
local WITH_NETWORK="${2}"
|
||||
local UPLINK="${3}"
|
||||
local BRIDGE="${4}"
|
||||
local PUBLIC_BRIDGE="${5}"
|
||||
local ROLLBACK_DELAY="${6}"
|
||||
local MASTER ADDR_JSON IP PREFIX GW MIGRATE MIGRATE_SCRIPT
|
||||
|
||||
[[ ${WITH_NETWORK} -eq ${FLAGS_TRUE} ]] || { info "réseau : ignoré (--nonetwork)"; return 0; }
|
||||
|
||||
info "réseau"
|
||||
|
||||
# Bridge public : créé à vide, réservé pour un usage futur. Aucune adresse,
|
||||
# aucun esclave, pas référencé dans la config agent.
|
||||
kvm_ensure_bridge "${DRY_RUN}" "${PUBLIC_BRIDGE}"
|
||||
|
||||
MASTER=""
|
||||
[[ ${DRY_RUN} -ne ${FLAGS_TRUE} ]] && MASTER=$(ip -j link show dev "${UPLINK}" | jq -r '.[0].master // ""')
|
||||
if [[ "${MASTER}" == "${BRIDGE}" ]]
|
||||
then
|
||||
info " ${UPLINK} déjà esclave de ${BRIDGE} — migration ignorée"
|
||||
kvm_ensure_bridge "${DRY_RUN}" "${BRIDGE}"
|
||||
return 0
|
||||
fi
|
||||
|
||||
# IP, préfixe et gateway dérivés de l'uplink : rien en dur. Le filtre
|
||||
# inet/global évite de tomber sur une IPv6 ou une adresse de lien.
|
||||
if [[ ${DRY_RUN} -eq ${FLAGS_TRUE} ]]
|
||||
then
|
||||
IP="<ip(${UPLINK})>"; PREFIX="<prefixlen>"; GW="<gateway>"
|
||||
else
|
||||
ADDR_JSON=$(ip -j addr show dev "${UPLINK}")
|
||||
IP=$(echo "${ADDR_JSON}" | jq -r '[.[0].addr_info[] | select(.family=="inet" and .scope=="global")][0].local // ""')
|
||||
PREFIX=$(echo "${ADDR_JSON}" | jq -r '[.[0].addr_info[] | select(.family=="inet" and .scope=="global")][0].prefixlen // ""')
|
||||
GW=$(ip -j route show default dev "${UPLINK}" | jq -r '.[0].gateway // ""')
|
||||
|
||||
[[ -n "${IP}" && -n "${PREFIX}" ]] || die "aucune adresse IPv4 globale sur ${UPLINK}"
|
||||
[[ -n "${GW}" ]] || die "aucune route par défaut via ${UPLINK}"
|
||||
fi
|
||||
|
||||
info " ${UPLINK} : ${IP}/${PREFIX} gw ${GW} → ${BRIDGE}"
|
||||
|
||||
# Cette étape coupe le réseau de l'host si elle échoue à mi-parcours, sans
|
||||
# console de secours. Deux garde-fous :
|
||||
# - un rollback armé avant d'y toucher : rien n'étant écrit sur disque, un
|
||||
# reboot ramène la configuration d'origine. Désarmé si le test passe.
|
||||
# - la séquence tourne sous systemd et non dans la session SSH : une
|
||||
# coupure de SSH ne l'interrompt plus à mi-chemin.
|
||||
run "${DRY_RUN}" "systemctl stop two-net-rollback.timer 2>/dev/null || true"
|
||||
run "${DRY_RUN}" "systemd-run --collect --unit=two-net-rollback --on-active=${ROLLBACK_DELAY} systemctl reboot"
|
||||
|
||||
MIGRATE=$(cat <<EOF
|
||||
set -e
|
||||
ip link show dev '${BRIDGE}' >/dev/null 2>&1 || ip link add name '${BRIDGE}' type bridge
|
||||
ip link set dev '${BRIDGE}' type bridge stp_state 0
|
||||
ip link set up dev '${BRIDGE}'
|
||||
ip link set '${UPLINK}' master '${BRIDGE}'
|
||||
ip addr add '${IP}/${PREFIX}' dev '${BRIDGE}'
|
||||
ip route replace default via '${GW}' dev '${BRIDGE}'
|
||||
ip addr del '${IP}/${PREFIX}' dev '${UPLINK}'
|
||||
pkill dhclient || true
|
||||
EOF
|
||||
)
|
||||
|
||||
if [[ ${DRY_RUN} -eq ${FLAGS_TRUE} ]]
|
||||
then
|
||||
echo "# systemd-run --wait --collect --service-type=oneshot --unit=two-net-migrate /bin/bash <<'EOF'"
|
||||
echo "${MIGRATE}" | sed -e 's/^/# /'
|
||||
echo "# EOF"
|
||||
else
|
||||
MIGRATE_SCRIPT=$(mktemp /run/two-net-migrate.XXXXXX)
|
||||
printf '%s\n' "${MIGRATE}" > "${MIGRATE_SCRIPT}"
|
||||
systemctl reset-failed two-net-migrate.service 2>/dev/null || true
|
||||
systemd-run --wait --collect --service-type=oneshot --unit=two-net-migrate \
|
||||
/bin/bash "${MIGRATE_SCRIPT}" \
|
||||
|| die "migration réseau échouée — rollback armé dans ${ROLLBACK_DELAY}s (reboot)"
|
||||
rm -f "${MIGRATE_SCRIPT}"
|
||||
|
||||
# Vérification de connectivité avant de désarmer : seul critère qui
|
||||
# distingue un succès d'un host qu'on vient d'isoler.
|
||||
ping -c 2 -W 2 "${GW}" >/dev/null 2>&1 \
|
||||
|| die "gateway ${GW} injoignable après migration — rollback armé dans ${ROLLBACK_DELAY}s (reboot)"
|
||||
fi
|
||||
|
||||
run "${DRY_RUN}" "systemctl stop two-net-rollback.timer 2>/dev/null || true"
|
||||
info " rollback désarmé"
|
||||
}
|
||||
|
||||
# Point d'entrée appelé par deploy.sh --bootstrap.
|
||||
bootstrap_host () {
|
||||
local DRY_RUN="${1}"
|
||||
local WITH_PACKAGES="${2}"
|
||||
local WITH_NETWORK="${3}"
|
||||
local UPLINK="${4}"
|
||||
local BRIDGE="${5}"
|
||||
local PUBLIC_BRIDGE="${6}"
|
||||
local ROLLBACK_DELAY="${7}"
|
||||
|
||||
kvm_packages "${DRY_RUN}" "${WITH_PACKAGES}"
|
||||
kvm_kernel "${DRY_RUN}"
|
||||
kvm_network "${DRY_RUN}" "${WITH_NETWORK}" "${UPLINK}" "${BRIDGE}" "${PUBLIC_BRIDGE}" "${ROLLBACK_DELAY}"
|
||||
}
|
||||
383
scripts/deploy.sh
Normal file → Executable file
383
scripts/deploy.sh
Normal file → Executable file
|
|
@ -10,6 +10,14 @@ case "${unameOut}" in
|
|||
esac
|
||||
|
||||
SCRIPT_PATH="scripts/deploy.sh"
|
||||
UNIT_DIR="/etc/systemd/system"
|
||||
SCRIPTS_DIR="/opt/two/scripts"
|
||||
# Asset listant les sommes de contrôle des autres, au format sha256sum.
|
||||
MANIFEST_NAME="SHA256SUMS"
|
||||
|
||||
info () { echo "== ${1}"; }
|
||||
warn () { echo "!! ${1}" >&2; }
|
||||
die () { echo "!! ${1}" >&2; exit 1; }
|
||||
|
||||
exec_with_dry_run () {
|
||||
if [[ ${1} -eq ${FLAGS_TRUE} ]]; then
|
||||
|
|
@ -18,7 +26,7 @@ exec_with_dry_run () {
|
|||
eval "${2}" 2> /tmp/error || \
|
||||
{
|
||||
echo -e "failed with following error";
|
||||
output=$(cat /tmp/error | sed -e "s/^/ error -> /g");
|
||||
local output; output=$(cat /tmp/error | sed -e "s/^/ error -> /g");
|
||||
echo -e "${output}";
|
||||
return 1;
|
||||
}
|
||||
|
|
@ -26,66 +34,357 @@ exec_with_dry_run () {
|
|||
return 0
|
||||
}
|
||||
|
||||
run () {
|
||||
exec_with_dry_run "${1}" "${2}" || die "échec : ${2}"
|
||||
}
|
||||
|
||||
check_latest_script () {
|
||||
REMOTE_URL="${1}"
|
||||
LOCAL_PATH="${2}"
|
||||
local REMOTE_URL="${1}"
|
||||
local LOCAL_PATH="${2}"
|
||||
local REMOTE_SUM LOCAL_SUM
|
||||
|
||||
REMOTE=$(curl --silent "${REMOTE_URL}" | sha256sum)
|
||||
LOCAL=$(cat ${LOCAL_PATH} | sha256sum)
|
||||
REMOTE_SUM=$(curl --silent "${REMOTE_URL}" | sha256sum)
|
||||
LOCAL_SUM=$(cat ${LOCAL_PATH} | sha256sum)
|
||||
|
||||
[[ "${REMOTE}" == "${LOCAL}" ]] || return 1
|
||||
[[ "${REMOTE_SUM}" == "${LOCAL_SUM}" ]] || return 1
|
||||
return 0
|
||||
}
|
||||
|
||||
download_binaries () {
|
||||
DRY_RUN="${1}"
|
||||
TAG="${2}"
|
||||
GIT_SERVER="${3}"
|
||||
REPO_PATH="${4}"
|
||||
list_active_units () {
|
||||
local PATTERN="${1}"
|
||||
systemctl list-units --state=active --no-legend --plain "${PATTERN}" 2>/dev/null \
|
||||
| awk '{print $1}' || true
|
||||
}
|
||||
|
||||
#'.[0].assets.[].browser_download_url'
|
||||
[[ "${TAG}" == "" ]] && TAG=$(curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/?limit=1" | jq -r '.[0].tag_name')
|
||||
echo "Deploy ${TAG} binaries"
|
||||
|
||||
BIN_PATH="/opt/two/${TAG}/bin/"
|
||||
LN_PATH="/opt/two/bin/"
|
||||
|
||||
exec_with_dry_run "${DRY_RUN}" "mkdir -p \"${BIN_PATH}\""
|
||||
exec_with_dry_run "${DRY_RUN}" "mkdir -p \"${LN_PATH}\""
|
||||
|
||||
curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/tags/${TAG}" | jq -c '.assets[]' | while read tmp
|
||||
# Units non instanciables du profil (agent.service) : celles qu'on arrête et
|
||||
# démarre nommément.
|
||||
profile_main_units () {
|
||||
local unit
|
||||
for unit in $(profile_units "${1}")
|
||||
do
|
||||
BINARY_NAME=$(echo "${tmp}" | jq -r '.name')
|
||||
BINARY_SHORT_NAME=$(echo "${BINARY_NAME}" | cut -d_ -f 1)
|
||||
BINARY_URL=$(echo "${tmp}" | jq -r '.browser_download_url')
|
||||
exec_with_dry_run "${DRY_RUN}" "curl --silent '${BINARY_URL}' -o '${BIN_PATH}${BINARY_NAME}'"
|
||||
exec_with_dry_run "${DRY_RUN}" "chmod +x '${BIN_PATH}${BINARY_NAME}'"
|
||||
exec_with_dry_run "${DRY_RUN}" "rm -f '${LN_PATH}${BINARY_SHORT_NAME}'"
|
||||
exec_with_dry_run "${DRY_RUN}" "ln -s '${BIN_PATH}${BINARY_NAME}' '${LN_PATH}${BINARY_SHORT_NAME}'"
|
||||
case "${unit}" in
|
||||
*@.service) continue ;;
|
||||
*) echo "${unit}" ;;
|
||||
esac
|
||||
done
|
||||
}
|
||||
|
||||
active_instances () {
|
||||
local unit
|
||||
for unit in $(profile_units "${1}")
|
||||
do
|
||||
case "${unit}" in
|
||||
*@.service) list_active_units "${unit%@.service}@*" ;;
|
||||
esac
|
||||
done
|
||||
}
|
||||
|
||||
stop_services () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local INSTANCES="${3}"
|
||||
local unit
|
||||
|
||||
for unit in $(profile_main_units "${PROFILE}")
|
||||
do
|
||||
run "${DRY_RUN}" "systemctl stop '${unit}'"
|
||||
done
|
||||
|
||||
for unit in ${INSTANCES}
|
||||
do
|
||||
run "${DRY_RUN}" "systemctl stop '${unit}'"
|
||||
done
|
||||
}
|
||||
|
||||
start_services () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local INSTANCES="${3}"
|
||||
local unit
|
||||
|
||||
for unit in ${INSTANCES}
|
||||
do
|
||||
exec_with_dry_run "${DRY_RUN}" "systemctl start '${unit}'" \
|
||||
|| warn "démarrage de ${unit} en échec — poursuite"
|
||||
done
|
||||
|
||||
for unit in $(profile_main_units "${PROFILE}")
|
||||
do
|
||||
run "${DRY_RUN}" "systemctl start '${unit}'"
|
||||
done
|
||||
}
|
||||
|
||||
profile_units () {
|
||||
case "${1}" in
|
||||
kvm) echo "agent.service dnsmasq@.service metadata@.service" ;;
|
||||
intel) echo "" ;;
|
||||
*) return 1 ;;
|
||||
esac
|
||||
}
|
||||
|
||||
profile_binaries () {
|
||||
case "${1}" in
|
||||
kvm) echo "agent metadata run-dnsmasq-in-netns.sh" ;;
|
||||
intel) echo "" ;;
|
||||
*) return 1 ;;
|
||||
esac
|
||||
}
|
||||
|
||||
profile_bootstrap () {
|
||||
case "${1}" in
|
||||
kvm) echo "bootstrap_kvm.sh" ;;
|
||||
intel) echo "" ;;
|
||||
*) return 1 ;;
|
||||
esac
|
||||
}
|
||||
|
||||
run_bootstrap () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local GIT_SERVER="${3}"
|
||||
local REPO_PATH="${4}"
|
||||
local BRANCH="${5}"
|
||||
local NAME BOOTSTRAP_URL
|
||||
|
||||
NAME=$(profile_bootstrap "${PROFILE}") || die "profil inconnu : ${PROFILE}"
|
||||
[[ -n "${NAME}" ]] || { info "bootstrap : rien à préparer pour le profil ${PROFILE}"; return 0; }
|
||||
|
||||
BOOTSTRAP_URL="${GIT_SERVER}${REPO_PATH}raw/branch/${BRANCH}scripts/${NAME}"
|
||||
info "bootstrap du profil ${PROFILE} (${SCRIPTS_DIR}/${NAME} ou ${BOOTSTRAP_URL})"
|
||||
|
||||
# shellcheck source=/dev/null
|
||||
[[ -f "${SCRIPTS_DIR}/${NAME}" ]] && . "${SCRIPTS_DIR}/${NAME}" || eval "$(curl --silent --fail "${BOOTSTRAP_URL}")"
|
||||
|
||||
command -v bootstrap_host >/dev/null 2>&1 \
|
||||
|| die "bootstrap_host() introuvable — ni ${SCRIPTS_DIR}/${NAME} ni ${BOOTSTRAP_URL} n'ont pu être chargés"
|
||||
|
||||
bootstrap_host "${DRY_RUN}" "${FLAGS_packages}" "${FLAGS_network}" \
|
||||
"${FLAGS_uplink}" "${FLAGS_bridge}" "${FLAGS_pub_bridge}" "${FLAGS_rollback}"
|
||||
}
|
||||
|
||||
resolve_tag () {
|
||||
local TAG="${1}"
|
||||
local GIT_SERVER="${2}"
|
||||
local REPO_PATH="${3}"
|
||||
|
||||
[[ -n "${TAG}" ]] && { echo "${TAG}"; return 0; }
|
||||
curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/?limit=1" | jq -r '.[0].tag_name'
|
||||
}
|
||||
|
||||
release_assets () {
|
||||
local GIT_SERVER="${1}"
|
||||
local REPO_PATH="${2}"
|
||||
local RELEASE_TAG="${3}"
|
||||
|
||||
curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/tags/${RELEASE_TAG}" | jq -c '.assets[]'
|
||||
}
|
||||
|
||||
asset_name () {
|
||||
echo "${1}" | jq -r --arg s "${2}" 'select(.name == $s or (.name | startswith($s + "_"))) | .name' | head -1
|
||||
}
|
||||
|
||||
asset_url () {
|
||||
echo "${1}" | jq -r --arg n "${2}" 'select(.name == $n) | .browser_download_url' | head -1
|
||||
}
|
||||
|
||||
release_manifest () {
|
||||
local ASSETS="${1}"
|
||||
local URL
|
||||
|
||||
URL=$(asset_url "${ASSETS}" "${MANIFEST_NAME}")
|
||||
[[ -n "${URL}" ]] || return 1
|
||||
curl --silent --fail "${URL}" || return 1
|
||||
}
|
||||
|
||||
manifest_sum () {
|
||||
awk -v n="${2}" '$2 == n { print $1 }' <<< "${1}" | head -1
|
||||
}
|
||||
|
||||
sum_matches () {
|
||||
local MANIFEST="${1}"
|
||||
local NAME="${2}"
|
||||
local FILE="${3}"
|
||||
local EXPECTED ACTUAL
|
||||
|
||||
EXPECTED=$(manifest_sum "${MANIFEST}" "${NAME}")
|
||||
[[ -n "${EXPECTED}" ]] || die "asset absent du manifeste ${MANIFEST_NAME} : ${NAME}"
|
||||
|
||||
[[ -f "${FILE}" ]] || return 1
|
||||
ACTUAL=$(sha256sum < "${FILE}" | cut -d' ' -f1)
|
||||
[[ "${ACTUAL}" == "${EXPECTED}" ]]
|
||||
}
|
||||
|
||||
fetch_asset () {
|
||||
local DRY_RUN="${1}"
|
||||
local MANIFEST="${2}"
|
||||
local ASSETS="${3}"
|
||||
local NAME="${4}"
|
||||
local DEST="${5}"
|
||||
local URL
|
||||
|
||||
if [[ -n "${MANIFEST}" ]] && sum_matches "${MANIFEST}" "${NAME}" "${DEST}"
|
||||
then
|
||||
info " ${NAME} déjà présent et conforme — téléchargement évité"
|
||||
return 0
|
||||
fi
|
||||
|
||||
URL=$(asset_url "${ASSETS}" "${NAME}")
|
||||
[[ -n "${URL}" ]] || die "asset absent de la release ${TAG} : ${NAME}"
|
||||
run "${DRY_RUN}" "curl --silent --fail '${URL}' -o '${DEST}'"
|
||||
|
||||
[[ ${DRY_RUN} -eq ${FLAGS_TRUE} ]] && return 0
|
||||
[[ -n "${MANIFEST}" ]] || return 0
|
||||
sum_matches "${MANIFEST}" "${NAME}" "${DEST}" \
|
||||
|| die "somme de contrôle incorrecte après téléchargement : ${NAME}"
|
||||
}
|
||||
|
||||
fetch_assets () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local ASSETS="${3}"
|
||||
local MANIFEST="${4}"
|
||||
local unit short FULL_NAME
|
||||
|
||||
info "assets de la release ${TAG} pour le profil ${PROFILE}"
|
||||
|
||||
run "${DRY_RUN}" "mkdir -p '${BIN_PATH}'"
|
||||
run "${DRY_RUN}" "mkdir -p '${UNIT_PATH}'"
|
||||
run "${DRY_RUN}" "mkdir -p '${LN_PATH}'"
|
||||
|
||||
for unit in $(profile_units "${PROFILE}")
|
||||
do
|
||||
fetch_asset "${DRY_RUN}" "${MANIFEST}" "${ASSETS}" "${unit}" "${UNIT_PATH}${unit}"
|
||||
done
|
||||
|
||||
for short in $(profile_binaries "${PROFILE}")
|
||||
do
|
||||
FULL_NAME=$(asset_name "${ASSETS}" "${short}")
|
||||
[[ -n "${FULL_NAME}" ]] || die "exécutable absent de la release ${TAG} : ${short}"
|
||||
fetch_asset "${DRY_RUN}" "${MANIFEST}" "${ASSETS}" "${FULL_NAME}" "${BIN_PATH}${FULL_NAME}"
|
||||
run "${DRY_RUN}" "chmod +x '${BIN_PATH}${FULL_NAME}'"
|
||||
done
|
||||
}
|
||||
|
||||
install_units () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local UNITS unit
|
||||
|
||||
UNITS=$(profile_units "${PROFILE}") || die "profil inconnu : ${PROFILE}"
|
||||
[[ -n "${UNITS}" ]] || { info "units : aucune pour le profil ${PROFILE}"; return 0; }
|
||||
|
||||
info "units du profil ${PROFILE}"
|
||||
|
||||
for unit in ${UNITS}
|
||||
do
|
||||
if [[ ${DRY_RUN} -ne ${FLAGS_TRUE} ]] && [[ ! -f "${UNIT_PATH}${unit}" ]]
|
||||
then
|
||||
die "unit absente des assets de la release ${TAG} : ${unit}"
|
||||
fi
|
||||
run "${DRY_RUN}" "install -m 0644 '${UNIT_PATH}${unit}' '${UNIT_DIR}/${unit}'"
|
||||
done
|
||||
|
||||
run "${DRY_RUN}" "systemctl daemon-reload"
|
||||
|
||||
for unit in ${UNITS}
|
||||
do
|
||||
case "${unit}" in
|
||||
*@.service) continue ;;
|
||||
esac
|
||||
run "${DRY_RUN}" "systemctl enable '${unit}'"
|
||||
done
|
||||
}
|
||||
|
||||
switch_binaries () {
|
||||
local DRY_RUN="${1}"
|
||||
local PROFILE="${2}"
|
||||
local ASSETS="${3}"
|
||||
local BINARIES INSTANCES short FULL_NAME
|
||||
|
||||
BINARIES=$(profile_binaries "${PROFILE}")
|
||||
[[ -n "${BINARIES}" ]] || { info "bascule : aucun exécutable pour le profil ${PROFILE}"; return 0; }
|
||||
|
||||
info "bascule des binaires"
|
||||
|
||||
INSTANCES=$(active_instances "${PROFILE}")
|
||||
|
||||
stop_services "${DRY_RUN}" "${PROFILE}" "${INSTANCES}"
|
||||
|
||||
for short in ${BINARIES}
|
||||
do
|
||||
FULL_NAME=$(asset_name "${ASSETS}" "${short}")
|
||||
run "${DRY_RUN}" "rm -f '${LN_PATH}${short}'"
|
||||
run "${DRY_RUN}" "ln -s '${BIN_PATH}${FULL_NAME}' '${LN_PATH}${short}'"
|
||||
done
|
||||
|
||||
start_services "${DRY_RUN}" "${PROFILE}" "${INSTANCES}"
|
||||
}
|
||||
|
||||
main () {
|
||||
[[ -f ./libs/shflags ]] && . ./libs/shflags || eval "$(curl --silent https://git.g3e.fr/H6N/tools/raw/branch/main/libs/shflags)"
|
||||
|
||||
DEFINE_boolean 'dryrun' false 'Enable dry-run mode' 'd'
|
||||
DEFINE_boolean 'up_script' true 'Upgrade script' 's'
|
||||
DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g'
|
||||
DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r'
|
||||
DEFINE_string 'branch' 'main/' 'Branch name' 'b'
|
||||
DEFINE_string 'tag' '' 'Tag name' 't'
|
||||
DEFINE_boolean 'dryrun' false 'Enable dry-run mode' 'd'
|
||||
DEFINE_boolean 'up_script' true 'Upgrade script' 's'
|
||||
DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g'
|
||||
DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r'
|
||||
DEFINE_string 'branch' 'main/' 'Branch name' 'b'
|
||||
DEFINE_string 'tag' '' 'Tag name' 't'
|
||||
DEFINE_string 'profile' 'kvm' 'Host profile: kvm' 'p'
|
||||
DEFINE_boolean 'bootstrap' false 'Prepare the host (packages, kernel, network)' 'i'
|
||||
DEFINE_boolean 'packages' true 'Install profile packages during --bootstrap' 'k'
|
||||
DEFINE_boolean 'network' true 'Configure bridges during --bootstrap' 'n'
|
||||
DEFINE_string 'uplink' 'eno1' 'Physical uplink interface' 'u'
|
||||
DEFINE_string 'bridge' 'br-000000' 'Main bridge, uplink is enslaved to it' 'B'
|
||||
DEFINE_string 'pub_bridge' 'br-public' 'Reserved empty bridge' 'P'
|
||||
DEFINE_integer 'rollback' 120 'Rollback reboot delay in seconds' 'R'
|
||||
DEFINE_boolean 'verify' true 'Check assets against SHA256SUMS' 'V'
|
||||
|
||||
FLAGS "$@" || exit $?
|
||||
eval set -- "${FLAGS_ARGV}"
|
||||
|
||||
SCRIPT_URL="${FLAGS_git_server}${FLAGS_repo_path}raw/branch/${FLAGS_branch}${SCRIPT_PATH}"
|
||||
check_latest_script "${SCRIPT_URL}" "${0}" || (
|
||||
[[ ${FLAGS_up_script} -eq ${FLAGS_TRUE} ]] && \
|
||||
exec_with_dry_run "${FLAGS_dryrun}" "curl --silent '${SCRIPT_URL}' -o '${0}'"
|
||||
exit 1
|
||||
)
|
||||
profile_units "${FLAGS_profile}" >/dev/null || die "profil inconnu : ${FLAGS_profile}"
|
||||
|
||||
download_binaries "${FLAGS_dryrun}" "${FLAGS_tag}" "${FLAGS_git_server}" "${FLAGS_repo_path}"
|
||||
if [[ ${FLAGS_up_script} -eq ${FLAGS_TRUE} ]]
|
||||
then
|
||||
local SCRIPT_URL="${FLAGS_git_server}${FLAGS_repo_path}raw/branch/${FLAGS_branch}${SCRIPT_PATH}"
|
||||
if ! check_latest_script "${SCRIPT_URL}" "${0}"
|
||||
then
|
||||
run "${FLAGS_dryrun}" "curl --silent '${SCRIPT_URL}' -o '${0}'"
|
||||
die "script local différent de la branche ${FLAGS_branch} — mis à jour, relancer"
|
||||
fi
|
||||
fi
|
||||
|
||||
[[ ${FLAGS_bootstrap} -eq ${FLAGS_TRUE} ]] && \
|
||||
run_bootstrap "${FLAGS_dryrun}" "${FLAGS_profile}" "${FLAGS_git_server}" "${FLAGS_repo_path}" "${FLAGS_branch}"
|
||||
|
||||
local TAG BIN_PATH UNIT_PATH LN_PATH ASSETS MANIFEST
|
||||
|
||||
TAG=$(resolve_tag "${FLAGS_tag}" "${FLAGS_git_server}" "${FLAGS_repo_path}")
|
||||
[[ -n "${TAG}" && "${TAG}" != "null" ]] || die "impossible de déterminer la release à déployer"
|
||||
|
||||
BIN_PATH="/opt/two/${TAG}/bin/"
|
||||
UNIT_PATH="/opt/two/${TAG}/units/"
|
||||
LN_PATH="/opt/two/bin/"
|
||||
|
||||
ASSETS=$(release_assets "${FLAGS_git_server}" "${FLAGS_repo_path}" "${TAG}")
|
||||
[[ -n "${ASSETS}" ]] || die "aucun asset dans la release ${TAG}"
|
||||
|
||||
# Manifeste absent = release antérieure à sa mise en place, ou CI incomplète.
|
||||
# C'est bloquant : une vérification qui se désactive d'elle-même ne vérifie
|
||||
# rien. --noverify est la sortie explicite.
|
||||
MANIFEST=""
|
||||
if [[ ${FLAGS_verify} -eq ${FLAGS_TRUE} ]]
|
||||
then
|
||||
MANIFEST=$(release_manifest "${ASSETS}") \
|
||||
|| die "${MANIFEST_NAME} absent de la release ${TAG} — --noverify pour déployer sans vérification"
|
||||
else
|
||||
warn "vérification des sommes de contrôle désactivée (--noverify)"
|
||||
fi
|
||||
|
||||
fetch_assets "${FLAGS_dryrun}" "${FLAGS_profile}" "${ASSETS}" "${MANIFEST}"
|
||||
|
||||
install_units "${FLAGS_dryrun}" "${FLAGS_profile}"
|
||||
switch_binaries "${FLAGS_dryrun}" "${FLAGS_profile}" "${ASSETS}"
|
||||
}
|
||||
|
||||
[[ "${BASH_SOURCE[0]}" == "${0}" ]] && (main "$@" || exit 1)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue