Compare commits

...

16 commits

Author SHA1 Message Date
d042a1ab5b
Merge branch 'feature-37' 2026-08-16 22:54:44 +02:00
ba189cf9f5
f-37: scripts/ci: add manifest handle
Some checks failed
Pre Release Workflow / set-release-target (push) Successful in 2s
Pre Release Workflow / build (metadata, amd64, linux) (push) Has been cancelled
Pre Release Workflow / upload-assets (agent.service, systemd/agent.service) (push) Has been cancelled
Pre Release Workflow / upload-assets (dnsmasq@.service, systemd/dnsmasq@.service) (push) Has been cancelled
Pre Release Workflow / upload-assets (metadata@.service, systemd/metadata@.service) (push) Has been cancelled
Pre Release Workflow / upload-assets (run-dnsmasq-in-netns.sh, scripts/run-dnsmasq-in-netns.sh) (push) Has been cancelled
Pre Release Workflow / checksums (push) Has been cancelled
Pre Release Workflow / prerelease (push) Has been cancelled
Pre Release Workflow / build (agent, amd64, linux) (push) Has been cancelled
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 17:01:53 +02:00
ec33948b8c
f-37: scripts: change to locals
All checks were successful
Pre Release Workflow / prerelease (push) Successful in 11s
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m36s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m30s
Pre Release Workflow / upload-assets (agent.service, systemd/agent.service) (push) Successful in 7s
Pre Release Workflow / upload-assets (dnsmasq@.service, systemd/dnsmasq@.service) (push) Successful in 7s
Pre Release Workflow / upload-assets (metadata@.service, systemd/metadata@.service) (push) Successful in 8s
Pre Release Workflow / upload-assets (run-dnsmasq-in-netns.sh, scripts/run-dnsmasq-in-netns.sh) (push) Successful in 6s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 13:52:06 +02:00
64438d64e6
f-37: script: clean commentaries
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 13:51:55 +02:00
3d2e42d9f1
f-37: scripts: refacto deploy.sh
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 13:51:54 +02:00
8d04472ec1
f-37: prometheus: add vms
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m35s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 8s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 01:00:00 +02:00
351248d250
f-37: script: add stop start agent
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 2s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m35s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 00:28:00 +02:00
c89e3f64de
f-37: add systemd-run
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-14 00:27:25 +02:00
2150530ef2
f-37: periph: add final touch #37
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m35s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 7s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:55:53 +02:00
fef04b10fc
f-37: migration: add state change at boot #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:43:16 +02:00
aa6611249b
f-37: test: add test for state validation
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:39:00 +02:00
576166f3f1
f-37: code: use state and check state #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:39:00 +02:00
4b5793806b
f-37: code: add state lib in vm, subnet and vpc #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:32:00 +02:00
78363194bf
f-37: add test for error state #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:13:25 +02:00
fa72b6cd47
f-37: code: add state error in return nil #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:12:32 +02:00
3edc319969
f-37: add new lib for state handle #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:02:36 +02:00
32 changed files with 1629 additions and 165 deletions

View file

@ -42,27 +42,69 @@ jobs:
goarch: ${{ matrix.goarch }} goarch: ${{ matrix.goarch }}
binari: ${{ matrix.binaries }} binari: ${{ matrix.binaries }}
secrets: inherit 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 runs-on: docker
needs: [set-release-target] needs: [set-release-target]
strategy: strategy:
matrix: matrix:
script: include:
- run-dnsmasq-in-netns.sh - 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: steps:
- uses: actions/checkout@v3 - uses: actions/checkout@v3
- name: Move asset - name: Move asset
run: | run: |
mkdir -p "dist" mkdir -p "dist"
cp scripts/${{ matrix.script }} dist/ cp "${{ matrix.path }}" dist/
- name: Upload script - name: Upload asset
uses: actions/upload-artifact@v3 uses: actions/upload-artifact@v3
with: with:
name: ${{ matrix.script }}-${{ needs.set-release-target.outputs.release_cible }} name: ${{ matrix.name }}-${{ needs.set-release-target.outputs.release_cible }}
path: dist/${{ matrix.script }} 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: prerelease:
runs-on: docker 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 uses: ./.forgejo/workflows/release.yml
with: with:
tag: ${{ needs.set-release-target.outputs.release_cible }} tag: ${{ needs.set-release-target.outputs.release_cible }}

View file

@ -91,7 +91,7 @@ paths:
"404": "404":
$ref: "#/components/responses/NotFound" $ref: "#/components/responses/NotFound"
"409": "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: content:
application/json: application/json:
schema: schema:
@ -146,7 +146,7 @@ paths:
schema: schema:
$ref: "#/components/schemas/Error" $ref: "#/components/schemas/Error"
"422": "422":
description: Subnet not found or not in created state description: Subnet not found, or not in creating/running state
content: content:
application/json: application/json:
schema: schema:
@ -185,6 +185,12 @@ paths:
$ref: "#/components/schemas/VM" $ref: "#/components/schemas/VM"
"404": "404":
$ref: "#/components/responses/NotFound" $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": "500":
$ref: "#/components/responses/InternalError" $ref: "#/components/responses/InternalError"
@ -274,6 +280,12 @@ paths:
$ref: "#/components/schemas/Subnet" $ref: "#/components/schemas/Subnet"
"404": "404":
$ref: "#/components/responses/NotFound" $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": "500":
$ref: "#/components/responses/InternalError" $ref: "#/components/responses/InternalError"
@ -314,8 +326,8 @@ components:
example: vp-00001 example: vp-00001
state: state:
type: string type: string
enum: [creating, created, deleting, deleted] enum: [creating, running, error, deleting, deleted]
example: created example: running
cidr: cidr:
type: string type: string
example: "10.0.0.0/16" example: "10.0.0.0/16"
@ -373,8 +385,8 @@ components:
example: sn-00001 example: sn-00001
state: state:
type: string type: string
enum: [creating, created, deleting, deleted] enum: [creating, running, error, deleting, deleted]
example: created example: running
vpc: vpc:
type: string type: string
example: vpc1 example: vpc1
@ -472,8 +484,8 @@ components:
example: vm-00001 example: vm-00001
state: state:
type: string type: string
enum: [starting, started, stopping, stopped] enum: [creating, running, error, deleting, deleted]
example: started example: running
metadata_port: metadata_port:
type: string type: string
example: "80" example: "80"

View file

@ -8,6 +8,7 @@ import (
agentapi "git.g3e.fr/syonad/two/internal/api/agent" agentapi "git.g3e.fr/syonad/two/internal/api/agent"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/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" agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"git.g3e.fr/syonad/two/pkg/logger" "git.g3e.fr/syonad/two/pkg/logger"
@ -31,6 +32,13 @@ func main() {
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false) db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
defer db.Close() 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 := worker.New(cfg.Worker.BufferSize)
q.Start(cfg.Worker.Count) q.Start(cfg.Worker.Count)

View file

@ -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) { func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) {
cmd := dispatcher.DeleteSubnetCommand{Name: name} cmd := dispatcher.DeleteSubnetCommand{Name: name}
if err := s.dispatcher.Prepare(cmd); err != nil { 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()}) json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
return return
} }

View file

@ -28,7 +28,7 @@ func TestListSubnets_Empty(t *testing.T) {
func TestListSubnets_WithData(t *testing.T) { func TestListSubnets_WithData(t *testing.T) {
s, db := newTestServer(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/vpc", "vpc-1")
kv.AddInDB(db, "subnet/sn-2/state", "creating") kv.AddInDB(db, "subnet/sn-2/state", "creating")
kv.AddInDB(db, "subnet/sn-2/vpc", "vpc-1") 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) { func TestPostSubnet_Created(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-new", Name: "sn-new",
VPC: "vpc-1", VPC: "vpc-1",
@ -91,7 +91,7 @@ func TestPostSubnet_MissingFields(t *testing.T) {
func TestPostSubnet_IfaceTypeOptional(t *testing.T) { func TestPostSubnet_IfaceTypeOptional(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-opt", Name: "sn-opt",
VPC: "vpc-1", VPC: "vpc-1",
@ -126,8 +126,8 @@ func TestPostSubnet_VPCNotFound(t *testing.T) {
func TestPostSubnet_Duplicate(t *testing.T) { func TestPostSubnet_Duplicate(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
kv.AddInDB(db, "subnet/sn-exist/state", "created") kv.AddInDB(db, "subnet/sn-exist/state", "running")
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-exist", Name: "sn-exist",
VPC: "vpc-1", VPC: "vpc-1",
@ -163,7 +163,7 @@ func TestPostSubnet_VPCDeleting(t *testing.T) {
func TestPostSubnet_BridgeMode_Success(t *testing.T) { func TestPostSubnet_BridgeMode_Success(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-br", Name: "sn-br",
VPC: "vpc-1", VPC: "vpc-1",
@ -190,7 +190,7 @@ func TestPostSubnet_BridgeMode_Success(t *testing.T) {
func TestPostSubnet_UnknownMode(t *testing.T) { func TestPostSubnet_UnknownMode(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-1", Name: "sn-1",
VPC: "vpc-1", VPC: "vpc-1",
@ -219,7 +219,7 @@ func TestPostSubnet_InvalidBody(t *testing.T) {
func TestGetSubnet_Found(t *testing.T) { func TestGetSubnet_Found(t *testing.T) {
s, db := newTestServer(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/vpc", "vpc-1")
kv.AddInDB(db, "subnet/sn-1/cidr", "10.0.0.0/24") kv.AddInDB(db, "subnet/sn-1/cidr", "10.0.0.0/24")
kv.AddInDB(db, "subnet/sn-1/interface_ip", "10.0.0.1") 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 var result Subnet
json.NewDecoder(w.Body).Decode(&result) 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) t.Errorf("résultat inattendu : %+v", result)
} }
if result.VPC != "vpc-1" { if result.VPC != "vpc-1" {
@ -261,7 +261,7 @@ func TestGetSubnet_EmptyName(t *testing.T) {
func TestDeleteSubnet_Success(t *testing.T) { func TestDeleteSubnet_Success(t *testing.T) {
s, db := newTestServer(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) req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-del", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.SubnetByNameHandler(w, req) 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) { func TestDeleteSubnet_NotFound(t *testing.T) {
s, _ := newTestServer(t) s, _ := newTestServer(t)
req := httptest.NewRequest(http.MethodDelete, "/subnets/inexistant", nil) req := httptest.NewRequest(http.MethodDelete, "/subnets/inexistant", nil)
@ -287,7 +298,7 @@ func TestDeleteSubnet_NotFound(t *testing.T) {
func TestSubnetByName_InvalidMethod(t *testing.T) { func TestSubnetByName_InvalidMethod(t *testing.T) {
s, db := newTestServer(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) req := httptest.NewRequest(http.MethodPut, "/subnets/sn-1", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.SubnetByNameHandler(w, req) s.SubnetByNameHandler(w, req)

View file

@ -15,7 +15,7 @@ import (
func TestVmFromDB_SingleDisk(t *testing.T) { func TestVmFromDB_SingleDisk(t *testing.T) {
entries := map[string]string{ entries := map[string]string{
"vm/vm-1/state": "started", "vm/vm-1/state": "running",
"vm/vm-1/subnet": "sn-1", "vm/vm-1/subnet": "sn-1",
"vm/vm-1/ip": "10.0.0.5", "vm/vm-1/ip": "10.0.0.5",
"vm/vm-1/metadata_port": "1234", "vm/vm-1/metadata_port": "1234",
@ -37,7 +37,7 @@ func TestVmFromDB_SingleDisk(t *testing.T) {
func TestVmFromDB_MultiDisk(t *testing.T) { func TestVmFromDB_MultiDisk(t *testing.T) {
entries := map[string]string{ entries := map[string]string{
"vm/vm-2/state": "started", "vm/vm-2/state": "running",
"vm/vm-2/subnet": "sn-1", "vm/vm-2/subnet": "sn-1",
"vm/vm-2/ip": "10.0.0.6", "vm/vm-2/ip": "10.0.0.6",
"vm/vm-2/metadata_port": "1235", "vm/vm-2/metadata_port": "1235",
@ -62,7 +62,7 @@ func TestVmFromDB_MultiDisk(t *testing.T) {
func TestVmFromDB_SlotGap(t *testing.T) { func TestVmFromDB_SlotGap(t *testing.T) {
// sdb absent — sda et sdc seulement // sdb absent — sda et sdc seulement
entries := map[string]string{ entries := map[string]string{
"vm/vm-3/state": "started", "vm/vm-3/state": "running",
"vm/vm-3/subnet": "sn-1", "vm/vm-3/subnet": "sn-1",
"vm/vm-3/ip": "10.0.0.7", "vm/vm-3/ip": "10.0.0.7",
"vm/vm-3/metadata_port": "1236", "vm/vm-3/metadata_port": "1236",
@ -94,7 +94,7 @@ func TestVmFromDB_SlotGap(t *testing.T) {
func TestStartVM_MultiDisk(t *testing.T) { func TestStartVM_MultiDisk(t *testing.T) {
s, db := newTestServer(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/vpc", "vpc-1")
body, _ := json.Marshal(VMCreateRequest{ body, _ := json.Marshal(VMCreateRequest{
@ -128,7 +128,7 @@ func TestStartVM_MultiDisk(t *testing.T) {
func TestStartVM_StorageReturnedInResponse(t *testing.T) { func TestStartVM_StorageReturnedInResponse(t *testing.T) {
s, db := newTestServer(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/vpc", "vpc-1")
body, _ := json.Marshal(VMCreateRequest{ body, _ := json.Marshal(VMCreateRequest{

View file

@ -28,7 +28,7 @@ func TestListVpcs_Empty(t *testing.T) {
func TestListVpcs_WithData(t *testing.T) { func TestListVpcs_WithData(t *testing.T) {
s, db := newTestServer(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") kv.AddInDB(db, "vpc/v2/state", "creating")
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodGet, "/vpcs", nil)) 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) { func TestPostVpc_Duplicate(t *testing.T) {
s, db := newTestServer(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"}) body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-exist", CIDR: "10.0.0.0/16"})
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body))) 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) { func TestGetVpc_Found(t *testing.T) {
s, db := newTestServer(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) req := httptest.NewRequest(http.MethodGet, "/vpcs/vpc-1", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.VpcByNameHandler(w, req) s.VpcByNameHandler(w, req)
@ -135,7 +135,7 @@ func TestGetVpc_Found(t *testing.T) {
} }
var result VPC var result VPC
json.NewDecoder(w.Body).Decode(&result) 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) t.Errorf("résultat inattendu : %+v", result)
} }
} }
@ -162,7 +162,7 @@ func TestGetVpc_EmptyName(t *testing.T) {
func TestDeleteVpc_Success(t *testing.T) { func TestDeleteVpc_Success(t *testing.T) {
s, db := newTestServer(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) req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-del", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.VpcByNameHandler(w, req) 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) { func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-busy/state", "created") kv.AddInDB(db, "vpc/vpc-busy/state", "running")
kv.AddInDB(db, "subnet/sn-1/state", "created") kv.AddInDB(db, "subnet/sn-1/state", "running")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy")
req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-busy", nil) req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-busy", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
@ -201,7 +228,7 @@ func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) {
func TestVpcByName_InvalidMethod(t *testing.T) { func TestVpcByName_InvalidMethod(t *testing.T) {
s, db := newTestServer(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) req := httptest.NewRequest(http.MethodPut, "/vpcs/vpc-1", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.VpcByNameHandler(w, req) s.VpcByNameHandler(w, req)

View file

@ -6,6 +6,7 @@ import (
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/state"
"git.g3e.fr/syonad/two/pkg/worker" "git.g3e.fr/syonad/two/pkg/worker"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
@ -13,6 +14,7 @@ import (
type Command interface { type Command interface {
Prepare(db *badger.DB, cfg *configuration.Config) error Prepare(db *badger.DB, cfg *configuration.Config) error
Execute(db *badger.DB, cfg *configuration.Config) error Execute(db *badger.DB, cfg *configuration.Config) error
Key() string
} }
type Dispatcher struct { type Dispatcher struct {
@ -43,6 +45,10 @@ func (d *Dispatcher) Dispatch(cmd Command) {
} }
if err != nil { if err != nil {
d.logger.Error("command failed", append(attrs, "error", err)...) 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 { } else {
d.logger.Info("command done", attrs...) d.logger.Info("command done", attrs...)
} }

View file

@ -4,8 +4,10 @@ import (
"errors" "errors"
"sync" "sync"
"testing" "testing"
"time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/state"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
@ -61,3 +63,66 @@ func TestDispatcher_Dispatch_ExecuteErrorLogged(t *testing.T) {
d.Dispatch(cmd) d.Dispatch(cmd)
wg.Wait() // Execute s'est terminé — l'erreur est loggée, pas propagée 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)
}
}

View file

@ -25,10 +25,18 @@ func newTestDispatcher(t *testing.T) (*Dispatcher, *badger.DB) {
// mockCmd implémente Command sans aucune dépendance système. // mockCmd implémente Command sans aucune dépendance système.
type mockCmd struct { type mockCmd struct {
key string
prepareFn func(*badger.DB, *configuration.Config) error prepareFn func(*badger.DB, *configuration.Config) error
executeFn 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 { func (m mockCmd) Prepare(db *badger.DB, cfg *configuration.Config) error {
return m.prepareFn(db, cfg) return m.prepareFn(db, cfg)
} }

View file

@ -6,6 +6,7 @@ import (
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" 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/internal/subnet"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
@ -22,6 +23,8 @@ type CreateSubnetCommand struct {
DefaultRoute bool DefaultRoute bool
} }
func (c CreateSubnetCommand) Key() string { return "subnet/" + c.Name }
func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error { func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error {
if c.Mode == "" { if c.Mode == "" {
c.Mode = "vxlan" 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 { if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil {
return fmt.Errorf("subnet %q already exists", c.Name) 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 { if err != nil {
return fmt.Errorf("vpc %q not found", c.VPC) 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) return fmt.Errorf("vpc %q is %s", c.VPC, vpcState)
} }
localIface, ok := cfg.Interfaces[c.IfaceType] localIface, ok := cfg.Interfaces[c.IfaceType]
if !ok { if !ok {
localIface = cfg.DefaultInterface 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+"/vpc", c.VPC)
kv.AddInDB(db, "subnet/"+c.Name+"/mode", c.Mode) kv.AddInDB(db, "subnet/"+c.Name+"/mode", c.Mode)
kv.AddInDB(db, "subnet/"+c.Name+"/local_iface", localIface) 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 { func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
for { for {
state, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state") vpcState, err := state.Get(db, "vpc/"+c.VPC)
if err != nil { if err != nil {
return fmt.Errorf("vpc %q not found while waiting", c.VPC) return fmt.Errorf("vpc %q not found while waiting", c.VPC)
} }
if state == "created" { if vpcState == state.Running {
break break
} }
if vpcState != state.Creating {
return fmt.Errorf("vpc %q is %s, cannot create subnet %q", c.VPC, vpcState, c.Name)
}
select { select {
case <-timeout: 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): case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
} }
} }
@ -79,22 +85,28 @@ type DeleteSubnetCommand struct {
Name string Name string
} }
func (c DeleteSubnetCommand) Key() string { return "subnet/" + c.Name }
func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) error { 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 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 { func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
if err := subnet.DeleteSubnet(db, c.Name); err != nil { if err := subnet.DeleteSubnet(db, c.Name); err != nil {
return err return err
} }
state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state") current, err := state.Get(db, c.Key())
if err != nil { if err != nil {
return err return err
} }
if state == "deleted" { if current == state.Deleted {
kv.DeleteInDB(db, "subnet/"+c.Name) kv.DeleteInDB(db, "subnet/"+c.Name)
} }
return nil return nil

View file

@ -17,7 +17,7 @@ func testCfg() *configuration.Config {
func TestCreateSubnetCommand_Prepare_Success(t *testing.T) { func TestCreateSubnetCommand_Prepare_Success(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "inconnu", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_Duplicate(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
kv.AddInDB(db, "subnet/sn-exist/state", "created") kv.AddInDB(db, "subnet/sn-exist/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-exist", VPC: "vpc-1", VxlanID: 100, Name: "sn-exist", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "bridge", Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "bridge", Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "vlan", Name: "sn-1", VPC: "vpc-1", Mode: "vlan",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "running")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", 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) { func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) {
_, db := newTestDispatcher(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"} cmd := DeleteSubnetCommand{Name: "sn-del"}
if err := cmd.Prepare(db, nil); err != nil { if err := cmd.Prepare(db, nil); err != nil {
t.Fatalf("Prepare a échoué : %v", err) 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) { func TestDeleteSubnetCommand_Prepare_NotFound(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
cmd := DeleteSubnetCommand{Name: "sn-inexistant"} cmd := DeleteSubnetCommand{Name: "sn-inexistant"}

View file

@ -8,6 +8,7 @@ import (
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" 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/internal/vm"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
@ -30,22 +31,24 @@ type StartVMCommand struct {
SSHKey string SSHKey string
} }
func (c StartVMCommand) Key() string { return "vm/" + c.Name }
func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil { if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil {
return fmt.Errorf("vm %q already exists", c.Name) 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 { if err != nil {
return fmt.Errorf("subnet %q not found", c.Subnet) 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) return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState)
} }
port, err := allocateMetadataPort(db) port, err := allocateMetadataPort(db)
if err != nil { if err != nil {
return fmt.Errorf("allocate metadata port: %w", err) 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+"/subnet", c.Subnet)
kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP) kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP)
kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port)) 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 { func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
for { for {
state, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state") subnetState, err := state.Get(db, "subnet/"+c.Subnet)
if err != nil { if err != nil {
return fmt.Errorf("subnet %q not found while waiting", c.Subnet) return fmt.Errorf("subnet %q not found while waiting", c.Subnet)
} }
if state == "created" { if subnetState == state.Running {
break break
} }
if subnetState != state.Creating {
return fmt.Errorf("subnet %q is %s, cannot start vm %q", c.Subnet, subnetState, c.Name)
}
select { select {
case <-timeout: 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): case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
} }
} }
@ -111,22 +117,28 @@ type StopVMCommand struct {
Name string Name string
} }
func (c StopVMCommand) Key() string { return "vm/" + c.Name }
func (c StopVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { 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 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 { func (c StopVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
if err := vm.StopVM(db, c.Name, cfg); err != nil { if err := vm.StopVM(db, c.Name, cfg); err != nil {
return err return err
} }
state, err := kv.GetFromDB(db, "vm/"+c.Name+"/state") current, err := state.Get(db, c.Key())
if err != nil { if err != nil {
return err return err
} }
if state == "stopped" { if current == state.Deleted {
kv.DeleteInDB(db, "vm/"+c.Name) kv.DeleteInDB(db, "vm/"+c.Name)
} }
return nil return nil

View file

@ -10,7 +10,7 @@ import (
func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) { func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) {
_, db := newTestDispatcher(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") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
cmd := StartVMCommand{ cmd := StartVMCommand{
@ -34,7 +34,7 @@ func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) {
func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) { func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) {
_, db := newTestDispatcher(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") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
cmd := StartVMCommand{ cmd := StartVMCommand{
@ -66,7 +66,7 @@ func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) {
func TestStartVMCommand_Prepare_SlotGap(t *testing.T) { func TestStartVMCommand_Prepare_SlotGap(t *testing.T) {
_, db := newTestDispatcher(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") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
// sdb absent au boot — slot réservé pour hotplug // 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) { func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) {
_, db := newTestDispatcher(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") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
cmd := StartVMCommand{ 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) { func TestStartVMCommand_Prepare_Duplicate(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vm/vm-exist/state", "started") kv.AddInDB(db, "vm/vm-exist/state", "running")
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-exist", Name: "vm-exist",

View file

@ -7,6 +7,7 @@ import (
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" 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/internal/vpc"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
@ -17,6 +18,8 @@ type CreateVPCCommand struct {
CIDR string CIDR string
} }
func (c CreateVPCCommand) Key() string { return "vpc/" + c.Name }
func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil { if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil {
return fmt.Errorf("vpc %q already exists", c.Name) 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 { if err := kv.AddInDB(db, "vpc/"+c.Name+"/cidr", c.CIDR); err != nil {
return err 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 { func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error {
@ -38,10 +41,16 @@ type DeleteVPCCommand struct {
Name string Name string
} }
func (c DeleteVPCCommand) Key() string { return "vpc/" + c.Name }
func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { 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) 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/") entries, err := kv.ListByPrefix(db, "subnet/")
if err != nil { if err != nil {
return fmt.Errorf("failed to list subnets: %w", err) return fmt.Errorf("failed to list subnets: %w", err)
@ -51,12 +60,12 @@ func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error
continue continue
} }
subnetName := strings.Split(key, "/")[1] subnetName := strings.Split(key, "/")[1]
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") s, err := state.Get(db, "subnet/"+subnetName)
if err != nil || (state != "deleting" && state != "deleted") { 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 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 { 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 { for key, value := range entries {
if strings.HasSuffix(key, "/vpc") && value == c.Name { if strings.HasSuffix(key, "/vpc") && value == c.Name {
subnetName := strings.Split(key, "/")[1] subnetName := strings.Split(key, "/")[1]
state, _ := kv.GetFromDB(db, "subnet/"+subnetName+"/state") s, _ := state.Get(db, "subnet/"+subnetName)
if state == "deleting" { if s == state.Deleting {
pending = true pending = true
break 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 { if err := vpc.DeleteVPC(db, c.Name); err != nil {
return err return err
} }
state, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state") current, err := state.Get(db, c.Key())
if err != nil { if err != nil {
return err return err
} }
if state == "deleted" { if current == state.Deleted {
kv.DeleteInDB(db, "vpc/"+c.Name) kv.DeleteInDB(db, "vpc/"+c.Name)
} }
return nil return nil

View file

@ -32,7 +32,7 @@ func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) { func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) {
_, db := newTestDispatcher(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"} cmd := CreateVPCCommand{Name: "vpc-exist", CIDR: "10.0.0.0/16"}
if err := cmd.Prepare(db, nil); err == nil { if err := cmd.Prepare(db, nil); err == nil {
t.Error("Prepare devrait échouer sur un VPC déjà existant") 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) { func TestDeleteVPCCommand_Prepare_Success(t *testing.T) {
_, db := newTestDispatcher(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"} cmd := DeleteVPCCommand{Name: "vpc-del"}
if err := cmd.Prepare(db, nil); err != nil { if err := cmd.Prepare(db, nil); err != nil {
t.Fatalf("Prepare a échoué : %v", err) 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) { func TestDeleteVPCCommand_Prepare_BlockedByActiveSubnet(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-busy/state", "created") kv.AddInDB(db, "vpc/vpc-busy/state", "running")
kv.AddInDB(db, "subnet/sn-1/state", "created") kv.AddInDB(db, "subnet/sn-1/state", "running")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy")
cmd := DeleteVPCCommand{Name: "vpc-busy"} cmd := DeleteVPCCommand{Name: "vpc-busy"}
if err := cmd.Prepare(db, nil); err == nil { 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) { func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleted(t *testing.T) {
_, db := newTestDispatcher(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/state", "deleted")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok")
cmd := DeleteVPCCommand{Name: "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) { func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleting(t *testing.T) {
_, db := newTestDispatcher(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/state", "deleting")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok")
cmd := DeleteVPCCommand{Name: "vpc-ok"} cmd := DeleteVPCCommand{Name: "vpc-ok"}

View 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 "", ""
}

View 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)
}
}

View file

@ -3,12 +3,13 @@ package agentmetrics
import ( import (
"strings" "strings"
"git.g3e.fr/syonad/two/internal/state"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
"github.com/prometheus/client_golang/prometheus" "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 // AgentCollector implements prometheus.Collector and exposes agent metrics
// by querying the BadgerDB on each scrape. // by querying the BadgerDB on each scrape.
@ -16,6 +17,7 @@ type AgentCollector struct {
db *badger.DB db *badger.DB
vpcsTotal *prometheus.Desc vpcsTotal *prometheus.Desc
subnetsTotal *prometheus.Desc subnetsTotal *prometheus.Desc
vmsTotal *prometheus.Desc
} }
func NewAgentCollector(db *badger.DB) *AgentCollector { func NewAgentCollector(db *badger.DB) *AgentCollector {
@ -31,23 +33,30 @@ func NewAgentCollector(db *badger.DB) *AgentCollector {
"Number of subnets by state.", "Number of subnets by state.",
[]string{"state"}, nil, []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) { func (c *AgentCollector) Describe(ch chan<- *prometheus.Desc) {
ch <- c.vpcsTotal ch <- c.vpcsTotal
ch <- c.subnetsTotal ch <- c.subnetsTotal
ch <- c.vmsTotal
} }
func (c *AgentCollector) Collect(ch chan<- prometheus.Metric) { func (c *AgentCollector) Collect(ch chan<- prometheus.Metric) {
c.collectStates(ch, "vpc/", c.vpcsTotal) c.collectStates(ch, "vpc/", c.vpcsTotal)
c.collectStates(ch, "subnet/", c.subnetsTotal) c.collectStates(ch, "subnet/", c.subnetsTotal)
c.collectStates(ch, "vm/", c.vmsTotal)
} }
// collectStates counts resources under the given DB prefix by their state value // collectStates counts resources under the given DB prefix by their state value
// and emits one gauge per state label. // and emits one gauge per state label.
func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix string, desc *prometheus.Desc) { 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 { for _, s := range allStates {
counts[s] = 0 counts[s] = 0
} }
@ -55,13 +64,18 @@ func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix strin
items, err := kv.ListByPrefix(c.db, prefix) items, err := kv.ListByPrefix(c.db, prefix)
if err == nil { if err == nil {
for key, val := range items { for key, val := range items {
if strings.HasSuffix(key, "/state") { if !strings.HasSuffix(key, "/state") {
counts[val]++ 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 { for _, s := range allStates {
ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[state], state) ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[s], string(s))
} }
} }

View 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)
}
}

View file

@ -1,5 +1,9 @@
package qemu package qemu
func ScopeName(vmName string) string {
return "two-vm-" + vmName + ".scope"
}
type DiskConfig struct { type DiskConfig struct {
Path string Path string
Dev string Dev string

View file

@ -103,9 +103,16 @@ func Start(cfg Config) error {
"-daemonize", "-daemonize",
) )
cmd := exec.Command("qemu-system-x86_64", args...) scopeArgs := append([]string{
if err := cmd.Run(); err != nil { "--scope",
return fmt.Errorf("qemu-system-x86_64: %w", err) "--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 return nil
} }

68
internal/state/state.go Normal file
View 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))
}

View 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")
}
}

View file

@ -7,18 +7,18 @@ import (
"git.g3e.fr/syonad/two/internal/ebtables" "git.g3e.fr/syonad/two/internal/ebtables"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "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" "git.g3e.fr/syonad/two/pkg/systemd"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
func CreateSubnet(db *badger.DB, subnetName string) error { 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 { if err != nil {
return err return err
} }
if state != "creating" { if current != state.Creating {
return nil return nil
} }
@ -31,7 +31,7 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
return err 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 { func createSubnet(db *badger.DB, subnetName string, d subnetData) error {

View file

@ -7,6 +7,7 @@ import (
"git.g3e.fr/syonad/two/internal/ebtables" "git.g3e.fr/syonad/two/internal/ebtables"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "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/db/kv"
"git.g3e.fr/syonad/two/pkg/systemd" "git.g3e.fr/syonad/two/pkg/systemd"
@ -14,11 +15,11 @@ import (
) )
func DeleteSubnet(db *badger.DB, subnetName string) error { 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 { if err != nil {
return err return err
} }
if state != "deleting" { if current != state.Deleting {
return nil return nil
} }
@ -44,7 +45,7 @@ func DeleteSubnet(db *badger.DB, subnetName string) error {
return fmt.Errorf("unknown subnet mode %q", d.mode) 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 { func stopDHCP(db *badger.DB, subnetName string, d subnetData) error {

View file

@ -12,17 +12,17 @@ import (
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/internal/qemu" "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" "github.com/dgraph-io/badger/v4"
) )
func StartVM(db *badger.DB, name string, cfg *configuration.Config) error { 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 { if err != nil {
return err return err
} }
if state != "starting" { if current != state.Creating {
return nil 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 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 { func copyFile(src, dst string) error {

View file

@ -12,17 +12,17 @@ import (
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/internal/qmp" "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" "github.com/dgraph-io/badger/v4"
) )
func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { 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 { if err != nil {
return err return err
} }
if state != "stopping" { if current != state.Deleting {
return nil return nil
} }
@ -65,7 +65,7 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
os.Remove(varsPath) 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) { func waitQMPDead(socketPath string, timeout, poll time.Duration) {

View file

@ -5,15 +5,15 @@ import (
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "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" "github.com/dgraph-io/badger/v4"
) )
func CreateVPC(db *badger.DB, name string) error { 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 return err
} else if state == "creating" { } else if current == state.Creating {
vpcID := strings.SplitN(name, "-", 2)[1] vpcID := strings.SplitN(name, "-", 2)[1]
if err := netns.Create(name); err != nil { if err := netns.Create(name); err != nil {
@ -48,7 +48,7 @@ func CreateVPC(db *badger.DB, name string) error {
}); err != nil { }); err != nil {
return err return err
} }
kv.AddInDB(db, "vpc/"+name+"/state", "created") return state.Set(db, "vpc/"+name, state.Running)
} }
return nil return nil
} }

View file

@ -5,15 +5,15 @@ import (
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "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" "github.com/dgraph-io/badger/v4"
) )
func DeleteVPC(db *badger.DB, name string) error { 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 return err
} else if state == "deleting" { } else if current == state.Deleting {
vpcID := strings.SplitN(name, "-", 2)[1] vpcID := strings.SplitN(name, "-", 2)[1]
if err := netif.DeleteLink("vp-" + vpcID + "-e"); err != nil { 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 { if err := netns.Delete(name); err != nil {
return err return err
} }
kv.AddInDB(db, "vpc/"+name+"/state", "deleted") return state.Set(db, "vpc/"+name, state.Deleted)
} }
return nil return nil

169
scripts/bootstrap_kvm.sh Executable file
View 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}"
}

385
scripts/deploy.sh Normal file → Executable file
View file

@ -10,6 +10,14 @@ case "${unameOut}" in
esac esac
SCRIPT_PATH="scripts/deploy.sh" 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 () { exec_with_dry_run () {
if [[ ${1} -eq ${FLAGS_TRUE} ]]; then if [[ ${1} -eq ${FLAGS_TRUE} ]]; then
@ -18,7 +26,7 @@ exec_with_dry_run () {
eval "${2}" 2> /tmp/error || \ eval "${2}" 2> /tmp/error || \
{ {
echo -e "failed with following 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}"; echo -e "${output}";
return 1; return 1;
} }
@ -26,67 +34,358 @@ exec_with_dry_run () {
return 0 return 0
} }
run () {
exec_with_dry_run "${1}" "${2}" || die "échec : ${2}"
}
check_latest_script () { check_latest_script () {
REMOTE_URL="${1}" local REMOTE_URL="${1}"
LOCAL_PATH="${2}" local LOCAL_PATH="${2}"
local REMOTE_SUM LOCAL_SUM
REMOTE=$(curl --silent "${REMOTE_URL}" | sha256sum) REMOTE_SUM=$(curl --silent "${REMOTE_URL}" | sha256sum)
LOCAL=$(cat ${LOCAL_PATH} | sha256sum) LOCAL_SUM=$(cat ${LOCAL_PATH} | sha256sum)
[[ "${REMOTE}" == "${LOCAL}" ]] || return 1 [[ "${REMOTE_SUM}" == "${LOCAL_SUM}" ]] || return 1
return 0 return 0
} }
download_binaries () { list_active_units () {
DRY_RUN="${1}" local PATTERN="${1}"
TAG="${2}" systemctl list-units --state=active --no-legend --plain "${PATTERN}" 2>/dev/null \
GIT_SERVER="${3}" | awk '{print $1}' || true
REPO_PATH="${4}" }
#'.[0].assets.[].browser_download_url' # Units non instanciables du profil (agent.service) : celles qu'on arrête et
[[ "${TAG}" == "" ]] && TAG=$(curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/?limit=1" | jq -r '.[0].tag_name') # démarre nommément.
echo "Deploy ${TAG} binaries" profile_main_units () {
local unit
BIN_PATH="/opt/two/${TAG}/bin/" for unit in $(profile_units "${1}")
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
do do
BINARY_NAME=$(echo "${tmp}" | jq -r '.name') case "${unit}" in
BINARY_SHORT_NAME=$(echo "${BINARY_NAME}" | cut -d_ -f 1) *@.service) continue ;;
BINARY_URL=$(echo "${tmp}" | jq -r '.browser_download_url') *) echo "${unit}" ;;
exec_with_dry_run "${DRY_RUN}" "curl --silent '${BINARY_URL}' -o '${BIN_PATH}${BINARY_NAME}'" esac
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}'"
done 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 () { main () {
[[ -f ./libs/shflags ]] && . ./libs/shflags || eval "$(curl --silent https://git.g3e.fr/H6N/tools/raw/branch/main/libs/shflags)" [[ -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 'dryrun' false 'Enable dry-run mode' 'd'
DEFINE_boolean 'up_script' true 'Upgrade script' 's' DEFINE_boolean 'up_script' true 'Upgrade script' 's'
DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g' DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g'
DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r' DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r'
DEFINE_string 'branch' 'main/' 'Branch name' 'b' DEFINE_string 'branch' 'main/' 'Branch name' 'b'
DEFINE_string 'tag' '' 'Tag name' 't' 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 $? FLAGS "$@" || exit $?
eval set -- "${FLAGS_ARGV}" eval set -- "${FLAGS_ARGV}"
SCRIPT_URL="${FLAGS_git_server}${FLAGS_repo_path}raw/branch/${FLAGS_branch}${SCRIPT_PATH}" profile_units "${FLAGS_profile}" >/dev/null || die "profil inconnu : ${FLAGS_profile}"
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
)
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) [[ "${BASH_SOURCE[0]}" == "${0}" ]] && (main "$@" || exit 1)
[[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1) [[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1)