diff --git a/.forgejo/workflows/prerelease.yml b/.forgejo/workflows/prerelease.yml index 1869374..242b5de 100644 --- a/.forgejo/workflows/prerelease.yml +++ b/.forgejo/workflows/prerelease.yml @@ -42,27 +42,69 @@ jobs: goarch: ${{ matrix.goarch }} binari: ${{ matrix.binaries }} secrets: inherit - upload-scripts: + # Scripts et units systemd publiés comme assets de release : deploy.sh les + # installe depuis la release, plus depuis la branche. Ajouter une entrée ici + # suffit à livrer un nouveau fichier. + upload-assets: runs-on: docker needs: [set-release-target] strategy: matrix: - script: - - run-dnsmasq-in-netns.sh + include: + - path: scripts/run-dnsmasq-in-netns.sh + name: run-dnsmasq-in-netns.sh + - path: systemd/agent.service + name: agent.service + - path: systemd/dnsmasq@.service + name: dnsmasq@.service + - path: systemd/metadata@.service + name: metadata@.service steps: - uses: actions/checkout@v3 - name: Move asset run: | mkdir -p "dist" - cp scripts/${{ matrix.script }} dist/ - - name: Upload script + cp "${{ matrix.path }}" dist/ + - name: Upload asset uses: actions/upload-artifact@v3 with: - name: ${{ matrix.script }}-${{ needs.set-release-target.outputs.release_cible }} - path: dist/${{ matrix.script }} + name: ${{ matrix.name }}-${{ needs.set-release-target.outputs.release_cible }} + path: dist/${{ matrix.name }} + # Manifeste des sommes de contrôle de tous les artefacts, au format sha256sum. + # Dépend de tous les jobs qui produisent des artefacts : ajouter un producteur + # sans l'ajouter ici donnerait un manifeste incomplet, donc un déploiement qui + # refuse des assets légitimes. + checksums: + runs-on: docker + needs: [set-release-target, build, upload-assets] + steps: + - name: Download all artifacts + uses: actions/download-artifact@v3 + with: + path: artifacts/ + - name: Générer SHA256SUMS + run: | + mkdir -p dist + # download-artifact place chaque artefact dans son propre + # sous-répertoire : on aplatit pour que le manifeste porte les noms + # d'assets, tels que deploy.sh les demandera. + find artifacts/ -type f -exec cp {} dist/ \; + # Généré depuis dist/ pour que les noms soient nus, et hors de dist/ + # pour que le manifeste ne se liste pas lui-même. + ( cd dist && sha256sum * ) > SHA256SUMS + cat SHA256SUMS + - name: Upload SHA256SUMS + uses: actions/upload-artifact@v3 + with: + name: SHA256SUMS-${{ needs.set-release-target.outputs.release_cible }} + path: SHA256SUMS prerelease: runs-on: docker - needs: [set-release-target, build] + # Tous les producteurs d'artefacts sont dans les needs : sans ça, le job + # release peut appeler download-artifact avant la fin des uploads et publier + # une release incomplète (assets ou manifeste manquants de façon non + # déterministe). + needs: [set-release-target, build, upload-assets, checksums] uses: ./.forgejo/workflows/release.yml with: tag: ${{ needs.set-release-target.outputs.release_cible }} diff --git a/api/agent.yaml b/api/agent.yaml index 7e9c461..0861353 100644 --- a/api/agent.yaml +++ b/api/agent.yaml @@ -91,7 +91,7 @@ paths: "404": $ref: "#/components/responses/NotFound" "409": - description: VPC not in a deletable state + description: VPC not deletable — only running or error states can be deleted, and all its subnets must be deleted first content: application/json: schema: @@ -146,7 +146,7 @@ paths: schema: $ref: "#/components/schemas/Error" "422": - description: Subnet not found or not in created state + description: Subnet not found, or not in creating/running state content: application/json: schema: @@ -185,6 +185,12 @@ paths: $ref: "#/components/schemas/VM" "404": $ref: "#/components/responses/NotFound" + "409": + description: VM not stoppable — only running or error states can be stopped + content: + application/json: + schema: + $ref: "#/components/schemas/Error" "500": $ref: "#/components/responses/InternalError" @@ -274,6 +280,12 @@ paths: $ref: "#/components/schemas/Subnet" "404": $ref: "#/components/responses/NotFound" + "409": + description: Subnet not deletable — only running or error states can be deleted + content: + application/json: + schema: + $ref: "#/components/schemas/Error" "500": $ref: "#/components/responses/InternalError" @@ -314,8 +326,8 @@ components: example: vp-00001 state: type: string - enum: [creating, created, deleting, deleted] - example: created + enum: [creating, running, error, deleting, deleted] + example: running cidr: type: string example: "10.0.0.0/16" @@ -373,8 +385,8 @@ components: example: sn-00001 state: type: string - enum: [creating, created, deleting, deleted] - example: created + enum: [creating, running, error, deleting, deleted] + example: running vpc: type: string example: vpc1 @@ -472,8 +484,8 @@ components: example: vm-00001 state: type: string - enum: [starting, started, stopping, stopped] - example: started + enum: [creating, running, error, deleting, deleted] + example: running metadata_port: type: string example: "80" diff --git a/cmd/agent/main.go b/cmd/agent/main.go index 7083076..0425d4e 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -8,6 +8,7 @@ import ( agentapi "git.g3e.fr/syonad/two/internal/api/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent" dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent" + "git.g3e.fr/syonad/two/internal/migration" agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent" "git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/logger" @@ -31,6 +32,13 @@ func main() { db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false) defer db.Close() + // Avant tout démarrage de service : la DB peut porter l'ancien vocabulaire + // d'états, et des ressources transitoires orphelines d'un arrêt précédent. + if err := migration.MigrateStates(db, log.With(slog.String("component", "migration"))); err != nil { + log.Error("failed to migrate states", "error", err) + return + } + q := worker.New(cfg.Worker.BufferSize) q.Start(cfg.Worker.Count) diff --git a/internal/api/agent/subnet.go b/internal/api/agent/subnet.go index 36ee9d7..f0671f0 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -68,7 +68,13 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string) func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) { cmd := dispatcher.DeleteSubnetCommand{Name: name} if err := s.dispatcher.Prepare(cmd); err != nil { - w.WriteHeader(http.StatusNotFound) + // 404 si la ressource n'existe pas, 409 si elle existe mais n'est pas + // dans un état supprimable — même convention que /vpcs et /vms. + if _, dbErr := kv.GetFromDB(s.db, "subnet/"+name+"/state"); dbErr != nil { + w.WriteHeader(http.StatusNotFound) + } else { + w.WriteHeader(http.StatusConflict) + } json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()}) return } diff --git a/internal/api/agent/subnet_test.go b/internal/api/agent/subnet_test.go index 36e42f7..77ab5e4 100644 --- a/internal/api/agent/subnet_test.go +++ b/internal/api/agent/subnet_test.go @@ -28,7 +28,7 @@ func TestListSubnets_Empty(t *testing.T) { func TestListSubnets_WithData(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") kv.AddInDB(db, "subnet/sn-2/state", "creating") kv.AddInDB(db, "subnet/sn-2/vpc", "vpc-1") @@ -55,7 +55,7 @@ func TestListSubnets_InvalidMethod(t *testing.T) { func TestPostSubnet_Created(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := SubnetCreateRequest{ Name: "sn-new", VPC: "vpc-1", @@ -91,7 +91,7 @@ func TestPostSubnet_MissingFields(t *testing.T) { func TestPostSubnet_IfaceTypeOptional(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := SubnetCreateRequest{ Name: "sn-opt", VPC: "vpc-1", @@ -126,8 +126,8 @@ func TestPostSubnet_VPCNotFound(t *testing.T) { func TestPostSubnet_Duplicate(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") - kv.AddInDB(db, "subnet/sn-exist/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") + kv.AddInDB(db, "subnet/sn-exist/state", "running") req := SubnetCreateRequest{ Name: "sn-exist", VPC: "vpc-1", @@ -163,7 +163,7 @@ func TestPostSubnet_VPCDeleting(t *testing.T) { func TestPostSubnet_BridgeMode_Success(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := SubnetCreateRequest{ Name: "sn-br", VPC: "vpc-1", @@ -190,7 +190,7 @@ func TestPostSubnet_BridgeMode_Success(t *testing.T) { func TestPostSubnet_UnknownMode(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := SubnetCreateRequest{ Name: "sn-1", VPC: "vpc-1", @@ -219,7 +219,7 @@ func TestPostSubnet_InvalidBody(t *testing.T) { func TestGetSubnet_Found(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") kv.AddInDB(db, "subnet/sn-1/cidr", "10.0.0.0/24") kv.AddInDB(db, "subnet/sn-1/interface_ip", "10.0.0.1") @@ -231,7 +231,7 @@ func TestGetSubnet_Found(t *testing.T) { } var result Subnet json.NewDecoder(w.Body).Decode(&result) - if result.Name != "sn-1" || result.State != "created" { + if result.Name != "sn-1" || result.State != "running" { t.Errorf("résultat inattendu : %+v", result) } if result.VPC != "vpc-1" { @@ -261,7 +261,7 @@ func TestGetSubnet_EmptyName(t *testing.T) { func TestDeleteSubnet_Success(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-del/state", "created") + kv.AddInDB(db, "subnet/sn-del/state", "running") req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-del", nil) w := httptest.NewRecorder() s.SubnetByNameHandler(w, req) @@ -275,6 +275,17 @@ func TestDeleteSubnet_Success(t *testing.T) { } } +func TestDeleteSubnet_ConflictWhileCreating(t *testing.T) { + s, db := newTestServer(t) + kv.AddInDB(db, "subnet/sn-wip/state", "creating") + req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-wip", nil) + w := httptest.NewRecorder() + s.SubnetByNameHandler(w, req) + if w.Code != http.StatusConflict { + t.Errorf("attendu 409, obtenu %d: %s", w.Code, w.Body.String()) + } +} + func TestDeleteSubnet_NotFound(t *testing.T) { s, _ := newTestServer(t) req := httptest.NewRequest(http.MethodDelete, "/subnets/inexistant", nil) @@ -287,7 +298,7 @@ func TestDeleteSubnet_NotFound(t *testing.T) { func TestSubnetByName_InvalidMethod(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") req := httptest.NewRequest(http.MethodPut, "/subnets/sn-1", nil) w := httptest.NewRecorder() s.SubnetByNameHandler(w, req) diff --git a/internal/api/agent/vm_test.go b/internal/api/agent/vm_test.go index 3af5968..db8bbd5 100644 --- a/internal/api/agent/vm_test.go +++ b/internal/api/agent/vm_test.go @@ -15,7 +15,7 @@ import ( func TestVmFromDB_SingleDisk(t *testing.T) { entries := map[string]string{ - "vm/vm-1/state": "started", + "vm/vm-1/state": "running", "vm/vm-1/subnet": "sn-1", "vm/vm-1/ip": "10.0.0.5", "vm/vm-1/metadata_port": "1234", @@ -37,7 +37,7 @@ func TestVmFromDB_SingleDisk(t *testing.T) { func TestVmFromDB_MultiDisk(t *testing.T) { entries := map[string]string{ - "vm/vm-2/state": "started", + "vm/vm-2/state": "running", "vm/vm-2/subnet": "sn-1", "vm/vm-2/ip": "10.0.0.6", "vm/vm-2/metadata_port": "1235", @@ -62,7 +62,7 @@ func TestVmFromDB_MultiDisk(t *testing.T) { func TestVmFromDB_SlotGap(t *testing.T) { // sdb absent — sda et sdc seulement entries := map[string]string{ - "vm/vm-3/state": "started", + "vm/vm-3/state": "running", "vm/vm-3/subnet": "sn-1", "vm/vm-3/ip": "10.0.0.7", "vm/vm-3/metadata_port": "1236", @@ -94,7 +94,7 @@ func TestVmFromDB_SlotGap(t *testing.T) { func TestStartVM_MultiDisk(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") body, _ := json.Marshal(VMCreateRequest{ @@ -128,7 +128,7 @@ func TestStartVM_MultiDisk(t *testing.T) { func TestStartVM_StorageReturnedInResponse(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") body, _ := json.Marshal(VMCreateRequest{ diff --git a/internal/api/agent/vpc_test.go b/internal/api/agent/vpc_test.go index 3516305..b5ba517 100644 --- a/internal/api/agent/vpc_test.go +++ b/internal/api/agent/vpc_test.go @@ -28,7 +28,7 @@ func TestListVpcs_Empty(t *testing.T) { func TestListVpcs_WithData(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/v1/state", "created") + kv.AddInDB(db, "vpc/v1/state", "running") kv.AddInDB(db, "vpc/v2/state", "creating") w := httptest.NewRecorder() s.VpcsHandler(w, httptest.NewRequest(http.MethodGet, "/vpcs", nil)) @@ -104,7 +104,7 @@ func TestPostVpc_InvalidCIDR(t *testing.T) { func TestPostVpc_Duplicate(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-exist/state", "created") + kv.AddInDB(db, "vpc/vpc-exist/state", "running") body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-exist", CIDR: "10.0.0.0/16"}) w := httptest.NewRecorder() s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body))) @@ -126,7 +126,7 @@ func TestPostVpc_InvalidBody(t *testing.T) { func TestGetVpc_Found(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := httptest.NewRequest(http.MethodGet, "/vpcs/vpc-1", nil) w := httptest.NewRecorder() s.VpcByNameHandler(w, req) @@ -135,7 +135,7 @@ func TestGetVpc_Found(t *testing.T) { } var result VPC json.NewDecoder(w.Body).Decode(&result) - if result.Name != "vpc-1" || result.State != "created" { + if result.Name != "vpc-1" || result.State != "running" { t.Errorf("résultat inattendu : %+v", result) } } @@ -162,7 +162,7 @@ func TestGetVpc_EmptyName(t *testing.T) { func TestDeleteVpc_Success(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-del/state", "created") + kv.AddInDB(db, "vpc/vpc-del/state", "running") req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-del", nil) w := httptest.NewRecorder() s.VpcByNameHandler(w, req) @@ -186,10 +186,37 @@ func TestDeleteVpc_NotFound(t *testing.T) { } } +func TestDeleteVpc_ConflictWhileCreating(t *testing.T) { + s, db := newTestServer(t) + kv.AddInDB(db, "vpc/vpc-wip/state", "creating") + req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-wip", nil) + w := httptest.NewRecorder() + s.VpcByNameHandler(w, req) + if w.Code != http.StatusConflict { + t.Errorf("attendu 409, obtenu %d: %s", w.Code, w.Body.String()) + } +} + +func TestDeleteVpc_AllowedFromError(t *testing.T) { + s, db := newTestServer(t) + kv.AddInDB(db, "vpc/vpc-ko/state", "error") + req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-ko", nil) + w := httptest.NewRecorder() + s.VpcByNameHandler(w, req) + if w.Code != http.StatusAccepted { + t.Fatalf("attendu 202, obtenu %d: %s", w.Code, w.Body.String()) + } + var result VPC + json.NewDecoder(w.Body).Decode(&result) + if result.State != "deleting" { + t.Errorf("state attendu deleting, obtenu %q", result.State) + } +} + func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-busy/state", "created") - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "vpc/vpc-busy/state", "running") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-busy", nil) w := httptest.NewRecorder() @@ -201,7 +228,7 @@ func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) { func TestVpcByName_InvalidMethod(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") req := httptest.NewRequest(http.MethodPut, "/vpcs/vpc-1", nil) w := httptest.NewRecorder() s.VpcByNameHandler(w, req) diff --git a/internal/dispatcher/agent/dispatcher.go b/internal/dispatcher/agent/dispatcher.go index 2897c8a..e4129a2 100644 --- a/internal/dispatcher/agent/dispatcher.go +++ b/internal/dispatcher/agent/dispatcher.go @@ -6,6 +6,7 @@ import ( "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/pkg/worker" "github.com/dgraph-io/badger/v4" ) @@ -13,6 +14,7 @@ import ( type Command interface { Prepare(db *badger.DB, cfg *configuration.Config) error Execute(db *badger.DB, cfg *configuration.Config) error + Key() string } type Dispatcher struct { @@ -43,6 +45,10 @@ func (d *Dispatcher) Dispatch(cmd Command) { } if err != nil { d.logger.Error("command failed", append(attrs, "error", err)...) + if setErr := state.Set(d.db, cmd.Key(), state.Error); setErr != nil { + d.logger.Error("failed to mark resource as errored", + "command", cmdType, "key", cmd.Key(), "error", setErr) + } } else { d.logger.Info("command done", attrs...) } diff --git a/internal/dispatcher/agent/dispatcher_test.go b/internal/dispatcher/agent/dispatcher_test.go index b4881d9..4475a82 100644 --- a/internal/dispatcher/agent/dispatcher_test.go +++ b/internal/dispatcher/agent/dispatcher_test.go @@ -4,8 +4,10 @@ import ( "errors" "sync" "testing" + "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/state" "github.com/dgraph-io/badger/v4" ) @@ -61,3 +63,66 @@ func TestDispatcher_Dispatch_ExecuteErrorLogged(t *testing.T) { d.Dispatch(cmd) wg.Wait() // Execute s'est terminé — l'erreur est loggée, pas propagée } + +func TestDispatcher_Dispatch_ExecuteErrorSetsErrorState(t *testing.T) { + d, db := newTestDispatcher(t) + if err := state.Set(db, "vpc/vpc-1", state.Creating); err != nil { + t.Fatalf("préparation du test : %v", err) + } + + done := make(chan struct{}) + cmd := mockCmd{ + key: "vpc/vpc-1", + prepareFn: func(*badger.DB, *configuration.Config) error { return nil }, + executeFn: func(*badger.DB, *configuration.Config) error { + return errors.New("execute failed") + }, + } + d.Dispatch(cmd) + + // L'état est écrit après le retour d'Execute : on scrute la DB. + go func() { + defer close(done) + for { + if s, err := state.Get(db, "vpc/vpc-1"); err == nil && s == state.Error { + return + } + time.Sleep(5 * time.Millisecond) + } + }() + select { + case <-done: + case <-time.After(2 * time.Second): + s, _ := state.Get(db, "vpc/vpc-1") + t.Fatalf("la ressource devrait être en %q, obtenu %q", state.Error, s) + } +} + +func TestDispatcher_Dispatch_ExecuteSuccessKeepsState(t *testing.T) { + d, db := newTestDispatcher(t) + if err := state.Set(db, "vpc/vpc-1", state.Running); err != nil { + t.Fatalf("préparation du test : %v", err) + } + + var wg sync.WaitGroup + wg.Add(1) + cmd := mockCmd{ + key: "vpc/vpc-1", + prepareFn: func(*badger.DB, *configuration.Config) error { return nil }, + executeFn: func(*badger.DB, *configuration.Config) error { + defer wg.Done() + return nil + }, + } + d.Dispatch(cmd) + wg.Wait() + time.Sleep(50 * time.Millisecond) // laisse le temps d'une écriture parasite + + s, err := state.Get(db, "vpc/vpc-1") + if err != nil { + t.Fatalf("Get a échoué : %v", err) + } + if s != state.Running { + t.Errorf("l'état devrait rester %q, obtenu %q", state.Running, s) + } +} diff --git a/internal/dispatcher/agent/helpers_test.go b/internal/dispatcher/agent/helpers_test.go index 2cee6fc..c14f6a3 100644 --- a/internal/dispatcher/agent/helpers_test.go +++ b/internal/dispatcher/agent/helpers_test.go @@ -25,10 +25,18 @@ func newTestDispatcher(t *testing.T) (*Dispatcher, *badger.DB) { // mockCmd implémente Command sans aucune dépendance système. type mockCmd struct { + key string prepareFn func(*badger.DB, *configuration.Config) error executeFn func(*badger.DB, *configuration.Config) error } +func (m mockCmd) Key() string { + if m.key == "" { + return "vpc/mock" + } + return m.key +} + func (m mockCmd) Prepare(db *badger.DB, cfg *configuration.Config) error { return m.prepareFn(db, cfg) } diff --git a/internal/dispatcher/agent/subnet_commands.go b/internal/dispatcher/agent/subnet_commands.go index e488842..6033a40 100644 --- a/internal/dispatcher/agent/subnet_commands.go +++ b/internal/dispatcher/agent/subnet_commands.go @@ -6,6 +6,7 @@ import ( "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/internal/subnet" "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" @@ -22,6 +23,8 @@ type CreateSubnetCommand struct { DefaultRoute bool } +func (c CreateSubnetCommand) Key() string { return "subnet/" + c.Name } + func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error { if c.Mode == "" { c.Mode = "vxlan" @@ -32,18 +35,18 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil { return fmt.Errorf("subnet %q already exists", c.Name) } - vpcState, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state") + vpcState, err := state.Get(db, "vpc/"+c.VPC) if err != nil { return fmt.Errorf("vpc %q not found", c.VPC) } - if vpcState == "deleting" || vpcState == "deleted" { + if vpcState != state.Creating && vpcState != state.Running { return fmt.Errorf("vpc %q is %s", c.VPC, vpcState) } localIface, ok := cfg.Interfaces[c.IfaceType] if !ok { localIface = cfg.DefaultInterface } - kv.AddInDB(db, "subnet/"+c.Name+"/state", "creating") + state.Set(db, c.Key(), state.Creating) kv.AddInDB(db, "subnet/"+c.Name+"/vpc", c.VPC) kv.AddInDB(db, "subnet/"+c.Name+"/mode", c.Mode) kv.AddInDB(db, "subnet/"+c.Name+"/local_iface", localIface) @@ -59,16 +62,19 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error { timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) for { - state, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state") + vpcState, err := state.Get(db, "vpc/"+c.VPC) if err != nil { return fmt.Errorf("vpc %q not found while waiting", c.VPC) } - if state == "created" { + if vpcState == state.Running { break } + if vpcState != state.Creating { + return fmt.Errorf("vpc %q is %s, cannot create subnet %q", c.VPC, vpcState, c.Name) + } select { case <-timeout: - return fmt.Errorf("timed out waiting for vpc %q to be created", c.VPC) + return fmt.Errorf("timed out waiting for vpc %q to be running", c.VPC) case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second): } } @@ -79,22 +85,28 @@ type DeleteSubnetCommand struct { Name string } +func (c DeleteSubnetCommand) Key() string { return "subnet/" + c.Name } + func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil { + current, err := state.Get(db, c.Key()) + if err != nil { return fmt.Errorf("subnet %q not found", c.Name) } - return kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting") + if !state.CanDelete(current) { + return fmt.Errorf("subnet %q cannot be deleted while %s", c.Name, current) + } + return state.Set(db, c.Key(), state.Deleting) } func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error { if err := subnet.DeleteSubnet(db, c.Name); err != nil { return err } - state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state") + current, err := state.Get(db, c.Key()) if err != nil { return err } - if state == "deleted" { + if current == state.Deleted { kv.DeleteInDB(db, "subnet/"+c.Name) } return nil diff --git a/internal/dispatcher/agent/subnet_commands_test.go b/internal/dispatcher/agent/subnet_commands_test.go index b6aacee..645f620 100644 --- a/internal/dispatcher/agent/subnet_commands_test.go +++ b/internal/dispatcher/agent/subnet_commands_test.go @@ -17,7 +17,7 @@ func testCfg() *configuration.Config { func TestCreateSubnetCommand_Prepare_Success(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -37,7 +37,7 @@ func TestCreateSubnetCommand_Prepare_Success(t *testing.T) { func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -51,7 +51,7 @@ func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) { func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "inconnu", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -65,8 +65,8 @@ func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing. func TestCreateSubnetCommand_Prepare_Duplicate(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") - kv.AddInDB(db, "subnet/sn-exist/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") + kv.AddInDB(db, "subnet/sn-exist/state", "running") cmd := CreateSubnetCommand{ Name: "sn-exist", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -113,7 +113,7 @@ func TestCreateSubnetCommand_Prepare_VPCDeleted(t *testing.T) { func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -130,7 +130,7 @@ func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) { func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", Mode: "bridge", IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -150,7 +150,7 @@ func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) { func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", Mode: "bridge", IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -163,7 +163,7 @@ func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) { func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", Mode: "vlan", IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -175,7 +175,7 @@ func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) { func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -195,7 +195,7 @@ func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) { func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/state", "running") cmd := CreateSubnetCommand{ Name: "sn-1", VPC: "vpc-1", VxlanID: 100, IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24", @@ -211,7 +211,7 @@ func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) { func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-del/state", "created") + kv.AddInDB(db, "subnet/sn-del/state", "running") cmd := DeleteSubnetCommand{Name: "sn-del"} if err := cmd.Prepare(db, nil); err != nil { t.Fatalf("Prepare a échoué : %v", err) @@ -222,6 +222,32 @@ func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) { } } +func TestDeleteSubnetCommand_Prepare_RefusedWhileCreating(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "subnet/sn-wip/state", "creating") + cmd := DeleteSubnetCommand{Name: "sn-wip"} + if err := cmd.Prepare(db, nil); err == nil { + t.Error("Prepare devrait refuser la suppression d'un subnet en creating") + } + s, _ := kv.GetFromDB(db, "subnet/sn-wip/state") + if s != "creating" { + t.Errorf("l'état ne devrait pas changer, obtenu %q", s) + } +} + +func TestDeleteSubnetCommand_Prepare_AllowedFromError(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "subnet/sn-ko/state", "error") + cmd := DeleteSubnetCommand{Name: "sn-ko"} + if err := cmd.Prepare(db, nil); err != nil { + t.Fatalf("Prepare devrait accepter un subnet en error : %v", err) + } + s, _ := kv.GetFromDB(db, "subnet/sn-ko/state") + if s != "deleting" { + t.Errorf("state attendu deleting, obtenu %q", s) + } +} + func TestDeleteSubnetCommand_Prepare_NotFound(t *testing.T) { _, db := newTestDispatcher(t) cmd := DeleteSubnetCommand{Name: "sn-inexistant"} diff --git a/internal/dispatcher/agent/vm_commands.go b/internal/dispatcher/agent/vm_commands.go index bdfeaec..6b50629 100644 --- a/internal/dispatcher/agent/vm_commands.go +++ b/internal/dispatcher/agent/vm_commands.go @@ -8,6 +8,7 @@ import ( "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/internal/vm" "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" @@ -30,22 +31,24 @@ type StartVMCommand struct { SSHKey string } +func (c StartVMCommand) Key() string { return "vm/" + c.Name } + func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil { return fmt.Errorf("vm %q already exists", c.Name) } - subnetState, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state") + subnetState, err := state.Get(db, "subnet/"+c.Subnet) if err != nil { return fmt.Errorf("subnet %q not found", c.Subnet) } - if subnetState == "deleting" || subnetState == "deleted" { + if subnetState != state.Creating && subnetState != state.Running { return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState) } port, err := allocateMetadataPort(db) if err != nil { return fmt.Errorf("allocate metadata port: %w", err) } - kv.AddInDB(db, "vm/"+c.Name+"/state", "starting") + state.Set(db, c.Key(), state.Creating) kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet) kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP) kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port)) @@ -91,16 +94,19 @@ func allocateMetadataPort(db *badger.DB) (int, error) { func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error { timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) for { - state, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state") + subnetState, err := state.Get(db, "subnet/"+c.Subnet) if err != nil { return fmt.Errorf("subnet %q not found while waiting", c.Subnet) } - if state == "created" { + if subnetState == state.Running { break } + if subnetState != state.Creating { + return fmt.Errorf("subnet %q is %s, cannot start vm %q", c.Subnet, subnetState, c.Name) + } select { case <-timeout: - return fmt.Errorf("timed out waiting for subnet %q to be created", c.Subnet) + return fmt.Errorf("timed out waiting for subnet %q to be running", c.Subnet) case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second): } } @@ -111,22 +117,28 @@ type StopVMCommand struct { Name string } +func (c StopVMCommand) Key() string { return "vm/" + c.Name } + func (c StopVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err != nil { + current, err := state.Get(db, c.Key()) + if err != nil { return fmt.Errorf("vm %q not found", c.Name) } - return kv.AddInDB(db, "vm/"+c.Name+"/state", "stopping") + if !state.CanDelete(current) { + return fmt.Errorf("vm %q cannot be stopped while %s", c.Name, current) + } + return state.Set(db, c.Key(), state.Deleting) } func (c StopVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error { if err := vm.StopVM(db, c.Name, cfg); err != nil { return err } - state, err := kv.GetFromDB(db, "vm/"+c.Name+"/state") + current, err := state.Get(db, c.Key()) if err != nil { return err } - if state == "stopped" { + if current == state.Deleted { kv.DeleteInDB(db, "vm/"+c.Name) } return nil diff --git a/internal/dispatcher/agent/vm_commands_test.go b/internal/dispatcher/agent/vm_commands_test.go index ba400b9..2747b0f 100644 --- a/internal/dispatcher/agent/vm_commands_test.go +++ b/internal/dispatcher/agent/vm_commands_test.go @@ -10,7 +10,7 @@ import ( func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") cmd := StartVMCommand{ @@ -34,7 +34,7 @@ func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) { func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") cmd := StartVMCommand{ @@ -66,7 +66,7 @@ func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) { func TestStartVMCommand_Prepare_SlotGap(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") // sdb absent au boot — slot réservé pour hotplug @@ -96,7 +96,7 @@ func TestStartVMCommand_Prepare_SlotGap(t *testing.T) { func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") cmd := StartVMCommand{ @@ -114,9 +114,54 @@ func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) { } } +// --- StopVMCommand.Prepare --- + +func TestStopVMCommand_Prepare_Success(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vm/vm-run/state", "running") + cmd := StopVMCommand{Name: "vm-run"} + if err := cmd.Prepare(db, nil); err != nil { + t.Fatalf("Prepare a échoué : %v", err) + } + s, _ := kv.GetFromDB(db, "vm/vm-run/state") + if s != "deleting" { + t.Errorf("state attendu deleting, obtenu %q", s) + } +} + +func TestStopVMCommand_Prepare_RefusedWhileCreating(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vm/vm-wip/state", "creating") + cmd := StopVMCommand{Name: "vm-wip"} + if err := cmd.Prepare(db, nil); err == nil { + t.Error("Prepare devrait refuser l'arrêt d'une VM en creating") + } + s, _ := kv.GetFromDB(db, "vm/vm-wip/state") + if s != "creating" { + t.Errorf("l'état ne devrait pas changer, obtenu %q", s) + } +} + +func TestStopVMCommand_Prepare_AllowedFromError(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vm/vm-ko/state", "error") + cmd := StopVMCommand{Name: "vm-ko"} + if err := cmd.Prepare(db, nil); err != nil { + t.Fatalf("Prepare devrait accepter une VM en error : %v", err) + } +} + +func TestStopVMCommand_Prepare_NotFound(t *testing.T) { + _, db := newTestDispatcher(t) + cmd := StopVMCommand{Name: "vm-inexistante"} + if err := cmd.Prepare(db, nil); err == nil { + t.Error("Prepare devrait échouer si la VM n'existe pas") + } +} + func TestStartVMCommand_Prepare_Duplicate(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vm/vm-exist/state", "started") + kv.AddInDB(db, "vm/vm-exist/state", "running") cmd := StartVMCommand{ Name: "vm-exist", diff --git a/internal/dispatcher/agent/vpc_commands.go b/internal/dispatcher/agent/vpc_commands.go index a129d88..9e26819 100644 --- a/internal/dispatcher/agent/vpc_commands.go +++ b/internal/dispatcher/agent/vpc_commands.go @@ -7,6 +7,7 @@ import ( "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/internal/vpc" "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" @@ -17,6 +18,8 @@ type CreateVPCCommand struct { CIDR string } +func (c CreateVPCCommand) Key() string { return "vpc/" + c.Name } + func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil { return fmt.Errorf("vpc %q already exists", c.Name) @@ -27,7 +30,7 @@ func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error if err := kv.AddInDB(db, "vpc/"+c.Name+"/cidr", c.CIDR); err != nil { return err } - return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating") + return state.Set(db, c.Key(), state.Creating) } func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { @@ -38,10 +41,16 @@ type DeleteVPCCommand struct { Name string } +func (c DeleteVPCCommand) Key() string { return "vpc/" + c.Name } + func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err != nil { + current, err := state.Get(db, c.Key()) + if err != nil { return fmt.Errorf("vpc %q not found", c.Name) } + if !state.CanDelete(current) { + return fmt.Errorf("vpc %q cannot be deleted while %s", c.Name, current) + } entries, err := kv.ListByPrefix(db, "subnet/") if err != nil { return fmt.Errorf("failed to list subnets: %w", err) @@ -51,12 +60,12 @@ func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error continue } subnetName := strings.Split(key, "/")[1] - state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") - if err != nil || (state != "deleting" && state != "deleted") { + s, err := state.Get(db, "subnet/"+subnetName) + if err != nil || (s != state.Deleting && s != state.Deleted) { return fmt.Errorf("subnet %q must be deleted before deleting vpc %q", subnetName, c.Name) } } - return kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting") + return state.Set(db, c.Key(), state.Deleting) } func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) error { @@ -70,8 +79,8 @@ func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) erro for key, value := range entries { if strings.HasSuffix(key, "/vpc") && value == c.Name { subnetName := strings.Split(key, "/")[1] - state, _ := kv.GetFromDB(db, "subnet/"+subnetName+"/state") - if state == "deleting" { + s, _ := state.Get(db, "subnet/"+subnetName) + if s == state.Deleting { pending = true break } @@ -89,11 +98,11 @@ func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) erro if err := vpc.DeleteVPC(db, c.Name); err != nil { return err } - state, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state") + current, err := state.Get(db, c.Key()) if err != nil { return err } - if state == "deleted" { + if current == state.Deleted { kv.DeleteInDB(db, "vpc/"+c.Name) } return nil diff --git a/internal/dispatcher/agent/vpc_commands_test.go b/internal/dispatcher/agent/vpc_commands_test.go index 9633733..78b3d98 100644 --- a/internal/dispatcher/agent/vpc_commands_test.go +++ b/internal/dispatcher/agent/vpc_commands_test.go @@ -32,7 +32,7 @@ func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) { func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-exist/state", "created") + kv.AddInDB(db, "vpc/vpc-exist/state", "running") cmd := CreateVPCCommand{Name: "vpc-exist", CIDR: "10.0.0.0/16"} if err := cmd.Prepare(db, nil); err == nil { t.Error("Prepare devrait échouer sur un VPC déjà existant") @@ -51,7 +51,7 @@ func TestCreateVPCCommand_Prepare_InvalidCIDR(t *testing.T) { func TestDeleteVPCCommand_Prepare_Success(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-del/state", "created") + kv.AddInDB(db, "vpc/vpc-del/state", "running") cmd := DeleteVPCCommand{Name: "vpc-del"} if err := cmd.Prepare(db, nil); err != nil { t.Fatalf("Prepare a échoué : %v", err) @@ -70,10 +70,45 @@ func TestDeleteVPCCommand_Prepare_NotFound(t *testing.T) { } } +func TestDeleteVPCCommand_Prepare_RefusedWhileCreating(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vpc/vpc-wip/state", "creating") + cmd := DeleteVPCCommand{Name: "vpc-wip"} + if err := cmd.Prepare(db, nil); err == nil { + t.Error("Prepare devrait refuser la suppression d'un VPC en creating") + } + s, _ := kv.GetFromDB(db, "vpc/vpc-wip/state") + if s != "creating" { + t.Errorf("l'état ne devrait pas changer, obtenu %q", s) + } +} + +func TestDeleteVPCCommand_Prepare_RefusedWhileDeleting(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vpc/vpc-gone/state", "deleting") + cmd := DeleteVPCCommand{Name: "vpc-gone"} + if err := cmd.Prepare(db, nil); err == nil { + t.Error("Prepare devrait refuser un VPC déjà en deleting") + } +} + +func TestDeleteVPCCommand_Prepare_AllowedFromError(t *testing.T) { + _, db := newTestDispatcher(t) + kv.AddInDB(db, "vpc/vpc-ko/state", "error") + cmd := DeleteVPCCommand{Name: "vpc-ko"} + if err := cmd.Prepare(db, nil); err != nil { + t.Fatalf("Prepare devrait accepter un VPC en error : %v", err) + } + s, _ := kv.GetFromDB(db, "vpc/vpc-ko/state") + if s != "deleting" { + t.Errorf("state attendu deleting, obtenu %q", s) + } +} + func TestDeleteVPCCommand_Prepare_BlockedByActiveSubnet(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-busy/state", "created") - kv.AddInDB(db, "subnet/sn-1/state", "created") + kv.AddInDB(db, "vpc/vpc-busy/state", "running") + kv.AddInDB(db, "subnet/sn-1/state", "running") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") cmd := DeleteVPCCommand{Name: "vpc-busy"} if err := cmd.Prepare(db, nil); err == nil { @@ -83,7 +118,7 @@ func TestDeleteVPCCommand_Prepare_BlockedByActiveSubnet(t *testing.T) { func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleted(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-ok/state", "created") + kv.AddInDB(db, "vpc/vpc-ok/state", "running") kv.AddInDB(db, "subnet/sn-1/state", "deleted") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok") cmd := DeleteVPCCommand{Name: "vpc-ok"} @@ -94,7 +129,7 @@ func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleted(t *testing.T) { func TestDeleteVPCCommand_Prepare_AllowedWhenSubnetDeleting(t *testing.T) { _, db := newTestDispatcher(t) - kv.AddInDB(db, "vpc/vpc-ok/state", "created") + kv.AddInDB(db, "vpc/vpc-ok/state", "running") kv.AddInDB(db, "subnet/sn-1/state", "deleting") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok") cmd := DeleteVPCCommand{Name: "vpc-ok"} diff --git a/internal/migration/state.go b/internal/migration/state.go new file mode 100644 index 0000000..99b5517 --- /dev/null +++ b/internal/migration/state.go @@ -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 "", "" +} diff --git a/internal/migration/state_test.go b/internal/migration/state_test.go new file mode 100644 index 0000000..85137a4 --- /dev/null +++ b/internal/migration/state_test.go @@ -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) + } +} diff --git a/internal/prometheus/agent/collector.go b/internal/prometheus/agent/collector.go index 04712ad..073dd4b 100644 --- a/internal/prometheus/agent/collector.go +++ b/internal/prometheus/agent/collector.go @@ -3,12 +3,13 @@ package agentmetrics import ( "strings" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" "github.com/prometheus/client_golang/prometheus" ) -var allStates = []string{"creating", "created", "deleting", "deleted"} +var allStates = state.All() // AgentCollector implements prometheus.Collector and exposes agent metrics // by querying the BadgerDB on each scrape. @@ -16,6 +17,7 @@ type AgentCollector struct { db *badger.DB vpcsTotal *prometheus.Desc subnetsTotal *prometheus.Desc + vmsTotal *prometheus.Desc } func NewAgentCollector(db *badger.DB) *AgentCollector { @@ -31,23 +33,30 @@ func NewAgentCollector(db *badger.DB) *AgentCollector { "Number of subnets by state.", []string{"state"}, nil, ), + vmsTotal: prometheus.NewDesc( + "syonad_vms_total", + "Number of VMs by state.", + []string{"state"}, nil, + ), } } func (c *AgentCollector) Describe(ch chan<- *prometheus.Desc) { ch <- c.vpcsTotal ch <- c.subnetsTotal + ch <- c.vmsTotal } func (c *AgentCollector) Collect(ch chan<- prometheus.Metric) { c.collectStates(ch, "vpc/", c.vpcsTotal) c.collectStates(ch, "subnet/", c.subnetsTotal) + c.collectStates(ch, "vm/", c.vmsTotal) } // collectStates counts resources under the given DB prefix by their state value // and emits one gauge per state label. func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix string, desc *prometheus.Desc) { - counts := make(map[string]float64, len(allStates)) + counts := make(map[state.State]float64, len(allStates)) for _, s := range allStates { counts[s] = 0 } @@ -55,13 +64,18 @@ func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix strin items, err := kv.ListByPrefix(c.db, prefix) if err == nil { for key, val := range items { - if strings.HasSuffix(key, "/state") { - counts[val]++ + if !strings.HasSuffix(key, "/state") { + continue + } + // Une valeur hors enum n'est pas comptée : la migration au + // démarrage les a toutes ramenées dans l'enum. + if s, err := state.Parse(val); err == nil { + counts[s]++ } } } - for _, state := range allStates { - ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[state], state) + for _, s := range allStates { + ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[s], string(s)) } } diff --git a/internal/prometheus/agent/collector_test.go b/internal/prometheus/agent/collector_test.go new file mode 100644 index 0000000..49272f9 --- /dev/null +++ b/internal/prometheus/agent/collector_test.go @@ -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//disk/ 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) + } +} diff --git a/internal/qemu/config.go b/internal/qemu/config.go index 9b50add..2177066 100644 --- a/internal/qemu/config.go +++ b/internal/qemu/config.go @@ -1,5 +1,9 @@ package qemu +func ScopeName(vmName string) string { + return "two-vm-" + vmName + ".scope" +} + type DiskConfig struct { Path string Dev string diff --git a/internal/qemu/start_linux.go b/internal/qemu/start_linux.go index 181bb1d..fa7bcf3 100644 --- a/internal/qemu/start_linux.go +++ b/internal/qemu/start_linux.go @@ -103,9 +103,16 @@ func Start(cfg Config) error { "-daemonize", ) - cmd := exec.Command("qemu-system-x86_64", args...) - if err := cmd.Run(); err != nil { - return fmt.Errorf("qemu-system-x86_64: %w", err) + scopeArgs := append([]string{ + "--scope", + "--unit=" + ScopeName(cfg.Name), + "--collect", + "qemu-system-x86_64", + }, args...) + + cmd := exec.Command("systemd-run", scopeArgs...) + if out, err := cmd.CombinedOutput(); err != nil { + return fmt.Errorf("systemd-run qemu-system-x86_64: %w: %s", err, strings.TrimSpace(string(out))) } return nil } diff --git a/internal/state/state.go b/internal/state/state.go new file mode 100644 index 0000000..c2fc9f4 --- /dev/null +++ b/internal/state/state.go @@ -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)) +} diff --git a/internal/state/state_test.go b/internal/state/state_test.go new file mode 100644 index 0000000..bcc0cf4 --- /dev/null +++ b/internal/state/state_test.go @@ -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") + } +} diff --git a/internal/subnet/create.go b/internal/subnet/create.go index a3faf46..75ef62e 100644 --- a/internal/subnet/create.go +++ b/internal/subnet/create.go @@ -7,18 +7,18 @@ import ( "git.g3e.fr/syonad/two/internal/ebtables" "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/pkg/systemd" "github.com/dgraph-io/badger/v4" ) func CreateSubnet(db *badger.DB, subnetName string) error { - state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") + current, err := state.Get(db, "subnet/"+subnetName) if err != nil { return err } - if state != "creating" { + if current != state.Creating { return nil } @@ -31,7 +31,7 @@ func CreateSubnet(db *badger.DB, subnetName string) error { return err } - return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created") + return state.Set(db, "subnet/"+subnetName, state.Running) } func createSubnet(db *badger.DB, subnetName string, d subnetData) error { diff --git a/internal/subnet/delete.go b/internal/subnet/delete.go index f745eda..1175fd5 100644 --- a/internal/subnet/delete.go +++ b/internal/subnet/delete.go @@ -7,6 +7,7 @@ import ( "git.g3e.fr/syonad/two/internal/ebtables" "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/systemd" @@ -14,11 +15,11 @@ import ( ) func DeleteSubnet(db *badger.DB, subnetName string) error { - state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") + current, err := state.Get(db, "subnet/"+subnetName) if err != nil { return err } - if state != "deleting" { + if current != state.Deleting { return nil } @@ -44,7 +45,7 @@ func DeleteSubnet(db *badger.DB, subnetName string) error { return fmt.Errorf("unknown subnet mode %q", d.mode) } - return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted") + return state.Set(db, "subnet/"+subnetName, state.Deleted) } func stopDHCP(db *badger.DB, subnetName string, d subnetData) error { diff --git a/internal/vm/create.go b/internal/vm/create.go index 3e0d4a7..74ed0bb 100644 --- a/internal/vm/create.go +++ b/internal/vm/create.go @@ -12,17 +12,17 @@ import ( "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/qemu" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/state" "github.com/dgraph-io/badger/v4" ) func StartVM(db *badger.DB, name string, cfg *configuration.Config) error { - state, err := kv.GetFromDB(db, "vm/"+name+"/state") + current, err := state.Get(db, "vm/"+name) if err != nil { return err } - if state != "starting" { + if current != state.Creating { return nil } @@ -84,7 +84,7 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error { return fmt.Errorf("start qemu: %w", err) } - return kv.AddInDB(db, "vm/"+name+"/state", "started") + return state.Set(db, "vm/"+name, state.Running) } func copyFile(src, dst string) error { diff --git a/internal/vm/delete.go b/internal/vm/delete.go index 3d808e6..f1490f0 100644 --- a/internal/vm/delete.go +++ b/internal/vm/delete.go @@ -12,17 +12,17 @@ import ( "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/qmp" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/state" "github.com/dgraph-io/badger/v4" ) func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { - state, err := kv.GetFromDB(db, "vm/"+name+"/state") + current, err := state.Get(db, "vm/"+name) if err != nil { return err } - if state != "stopping" { + if current != state.Deleting { return nil } @@ -65,7 +65,7 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { os.Remove(varsPath) } - return kv.AddInDB(db, "vm/"+name+"/state", "stopped") + return state.Set(db, "vm/"+name, state.Deleted) } func waitQMPDead(socketPath string, timeout, poll time.Duration) { diff --git a/internal/vpc/create.go b/internal/vpc/create.go index 148f70c..52e2567 100644 --- a/internal/vpc/create.go +++ b/internal/vpc/create.go @@ -5,15 +5,15 @@ import ( "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/state" "github.com/dgraph-io/badger/v4" ) func CreateVPC(db *badger.DB, name string) error { - if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil { + if current, err := state.Get(db, "vpc/"+name); err != nil { return err - } else if state == "creating" { + } else if current == state.Creating { vpcID := strings.SplitN(name, "-", 2)[1] if err := netns.Create(name); err != nil { @@ -48,7 +48,7 @@ func CreateVPC(db *badger.DB, name string) error { }); err != nil { return err } - kv.AddInDB(db, "vpc/"+name+"/state", "created") + return state.Set(db, "vpc/"+name, state.Running) } return nil } diff --git a/internal/vpc/delete.go b/internal/vpc/delete.go index dbd3a59..ff4d83f 100644 --- a/internal/vpc/delete.go +++ b/internal/vpc/delete.go @@ -5,15 +5,15 @@ import ( "git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netns" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/state" "github.com/dgraph-io/badger/v4" ) func DeleteVPC(db *badger.DB, name string) error { - if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil { + if current, err := state.Get(db, "vpc/"+name); err != nil { return err - } else if state == "deleting" { + } else if current == state.Deleting { vpcID := strings.SplitN(name, "-", 2)[1] if err := netif.DeleteLink("vp-" + vpcID + "-e"); err != nil { @@ -23,7 +23,7 @@ func DeleteVPC(db *badger.DB, name string) error { if err := netns.Delete(name); err != nil { return err } - kv.AddInDB(db, "vpc/"+name+"/state", "deleted") + return state.Set(db, "vpc/"+name, state.Deleted) } return nil diff --git a/scripts/bootstrap_kvm.sh b/scripts/bootstrap_kvm.sh new file mode 100755 index 0000000..e239b7c --- /dev/null +++ b/scripts/bootstrap_kvm.sh @@ -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_.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=""; PREFIX=""; GW="" + 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 </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}" +} diff --git a/scripts/deploy.sh b/scripts/deploy.sh old mode 100644 new mode 100755 index 7532b14..07d7668 --- a/scripts/deploy.sh +++ b/scripts/deploy.sh @@ -10,6 +10,14 @@ case "${unameOut}" in esac SCRIPT_PATH="scripts/deploy.sh" +UNIT_DIR="/etc/systemd/system" +SCRIPTS_DIR="/opt/two/scripts" +# Asset listant les sommes de contrôle des autres, au format sha256sum. +MANIFEST_NAME="SHA256SUMS" + +info () { echo "== ${1}"; } +warn () { echo "!! ${1}" >&2; } +die () { echo "!! ${1}" >&2; exit 1; } exec_with_dry_run () { if [[ ${1} -eq ${FLAGS_TRUE} ]]; then @@ -18,7 +26,7 @@ exec_with_dry_run () { eval "${2}" 2> /tmp/error || \ { echo -e "failed with following error"; - output=$(cat /tmp/error | sed -e "s/^/ error -> /g"); + local output; output=$(cat /tmp/error | sed -e "s/^/ error -> /g"); echo -e "${output}"; return 1; } @@ -26,67 +34,358 @@ exec_with_dry_run () { return 0 } +run () { + exec_with_dry_run "${1}" "${2}" || die "échec : ${2}" +} + check_latest_script () { - REMOTE_URL="${1}" - LOCAL_PATH="${2}" + local REMOTE_URL="${1}" + local LOCAL_PATH="${2}" + local REMOTE_SUM LOCAL_SUM - REMOTE=$(curl --silent "${REMOTE_URL}" | sha256sum) - LOCAL=$(cat ${LOCAL_PATH} | sha256sum) + REMOTE_SUM=$(curl --silent "${REMOTE_URL}" | sha256sum) + LOCAL_SUM=$(cat ${LOCAL_PATH} | sha256sum) - [[ "${REMOTE}" == "${LOCAL}" ]] || return 1 + [[ "${REMOTE_SUM}" == "${LOCAL_SUM}" ]] || return 1 return 0 } -download_binaries () { - DRY_RUN="${1}" - TAG="${2}" - GIT_SERVER="${3}" - REPO_PATH="${4}" +list_active_units () { + local PATTERN="${1}" + systemctl list-units --state=active --no-legend --plain "${PATTERN}" 2>/dev/null \ + | awk '{print $1}' || true +} - #'.[0].assets.[].browser_download_url' - [[ "${TAG}" == "" ]] && TAG=$(curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/?limit=1" | jq -r '.[0].tag_name') - echo "Deploy ${TAG} binaries" - - BIN_PATH="/opt/two/${TAG}/bin/" - LN_PATH="/opt/two/bin/" - - exec_with_dry_run "${DRY_RUN}" "mkdir -p \"${BIN_PATH}\"" - exec_with_dry_run "${DRY_RUN}" "mkdir -p \"${LN_PATH}\"" - - curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/tags/${TAG}" | jq -c '.assets[]' | while read tmp +# Units non instanciables du profil (agent.service) : celles qu'on arrête et +# démarre nommément. +profile_main_units () { + local unit + for unit in $(profile_units "${1}") do - BINARY_NAME=$(echo "${tmp}" | jq -r '.name') - BINARY_SHORT_NAME=$(echo "${BINARY_NAME}" | cut -d_ -f 1) - BINARY_URL=$(echo "${tmp}" | jq -r '.browser_download_url') - exec_with_dry_run "${DRY_RUN}" "curl --silent '${BINARY_URL}' -o '${BIN_PATH}${BINARY_NAME}'" - exec_with_dry_run "${DRY_RUN}" "chmod +x '${BIN_PATH}${BINARY_NAME}'" - exec_with_dry_run "${DRY_RUN}" "rm -f '${LN_PATH}${BINARY_SHORT_NAME}'" - exec_with_dry_run "${DRY_RUN}" "ln -s '${BIN_PATH}${BINARY_NAME}' '${LN_PATH}${BINARY_SHORT_NAME}'" + case "${unit}" in + *@.service) continue ;; + *) echo "${unit}" ;; + esac done } +active_instances () { + local unit + for unit in $(profile_units "${1}") + do + case "${unit}" in + *@.service) list_active_units "${unit%@.service}@*" ;; + esac + done +} + +stop_services () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local INSTANCES="${3}" + local unit + + for unit in $(profile_main_units "${PROFILE}") + do + run "${DRY_RUN}" "systemctl stop '${unit}'" + done + + for unit in ${INSTANCES} + do + run "${DRY_RUN}" "systemctl stop '${unit}'" + done +} + +start_services () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local INSTANCES="${3}" + local unit + + for unit in ${INSTANCES} + do + exec_with_dry_run "${DRY_RUN}" "systemctl start '${unit}'" \ + || warn "démarrage de ${unit} en échec — poursuite" + done + + for unit in $(profile_main_units "${PROFILE}") + do + run "${DRY_RUN}" "systemctl start '${unit}'" + done +} + +profile_units () { + case "${1}" in + kvm) echo "agent.service dnsmasq@.service metadata@.service" ;; + intel) echo "" ;; + *) return 1 ;; + esac +} + +profile_binaries () { + case "${1}" in + kvm) echo "agent metadata run-dnsmasq-in-netns.sh" ;; + intel) echo "" ;; + *) return 1 ;; + esac +} + +profile_bootstrap () { + case "${1}" in + kvm) echo "bootstrap_kvm.sh" ;; + intel) echo "" ;; + *) return 1 ;; + esac +} + +run_bootstrap () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local GIT_SERVER="${3}" + local REPO_PATH="${4}" + local BRANCH="${5}" + local NAME BOOTSTRAP_URL + + NAME=$(profile_bootstrap "${PROFILE}") || die "profil inconnu : ${PROFILE}" + [[ -n "${NAME}" ]] || { info "bootstrap : rien à préparer pour le profil ${PROFILE}"; return 0; } + + BOOTSTRAP_URL="${GIT_SERVER}${REPO_PATH}raw/branch/${BRANCH}scripts/${NAME}" + info "bootstrap du profil ${PROFILE} (${SCRIPTS_DIR}/${NAME} ou ${BOOTSTRAP_URL})" + + # shellcheck source=/dev/null + [[ -f "${SCRIPTS_DIR}/${NAME}" ]] && . "${SCRIPTS_DIR}/${NAME}" || eval "$(curl --silent --fail "${BOOTSTRAP_URL}")" + + command -v bootstrap_host >/dev/null 2>&1 \ + || die "bootstrap_host() introuvable — ni ${SCRIPTS_DIR}/${NAME} ni ${BOOTSTRAP_URL} n'ont pu être chargés" + + bootstrap_host "${DRY_RUN}" "${FLAGS_packages}" "${FLAGS_network}" \ + "${FLAGS_uplink}" "${FLAGS_bridge}" "${FLAGS_pub_bridge}" "${FLAGS_rollback}" +} + +resolve_tag () { + local TAG="${1}" + local GIT_SERVER="${2}" + local REPO_PATH="${3}" + + [[ -n "${TAG}" ]] && { echo "${TAG}"; return 0; } + curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/?limit=1" | jq -r '.[0].tag_name' +} + +release_assets () { + local GIT_SERVER="${1}" + local REPO_PATH="${2}" + local RELEASE_TAG="${3}" + + curl --silent "${GIT_SERVER}api/v1/repos/${REPO_PATH}releases/tags/${RELEASE_TAG}" | jq -c '.assets[]' +} + +asset_name () { + echo "${1}" | jq -r --arg s "${2}" 'select(.name == $s or (.name | startswith($s + "_"))) | .name' | head -1 +} + +asset_url () { + echo "${1}" | jq -r --arg n "${2}" 'select(.name == $n) | .browser_download_url' | head -1 +} + +release_manifest () { + local ASSETS="${1}" + local URL + + URL=$(asset_url "${ASSETS}" "${MANIFEST_NAME}") + [[ -n "${URL}" ]] || return 1 + curl --silent --fail "${URL}" || return 1 +} + +manifest_sum () { + awk -v n="${2}" '$2 == n { print $1 }' <<< "${1}" | head -1 +} + +sum_matches () { + local MANIFEST="${1}" + local NAME="${2}" + local FILE="${3}" + local EXPECTED ACTUAL + + EXPECTED=$(manifest_sum "${MANIFEST}" "${NAME}") + [[ -n "${EXPECTED}" ]] || die "asset absent du manifeste ${MANIFEST_NAME} : ${NAME}" + + [[ -f "${FILE}" ]] || return 1 + ACTUAL=$(sha256sum < "${FILE}" | cut -d' ' -f1) + [[ "${ACTUAL}" == "${EXPECTED}" ]] +} + +fetch_asset () { + local DRY_RUN="${1}" + local MANIFEST="${2}" + local ASSETS="${3}" + local NAME="${4}" + local DEST="${5}" + local URL + + if [[ -n "${MANIFEST}" ]] && sum_matches "${MANIFEST}" "${NAME}" "${DEST}" + then + info " ${NAME} déjà présent et conforme — téléchargement évité" + return 0 + fi + + URL=$(asset_url "${ASSETS}" "${NAME}") + [[ -n "${URL}" ]] || die "asset absent de la release ${TAG} : ${NAME}" + run "${DRY_RUN}" "curl --silent --fail '${URL}' -o '${DEST}'" + + [[ ${DRY_RUN} -eq ${FLAGS_TRUE} ]] && return 0 + [[ -n "${MANIFEST}" ]] || return 0 + sum_matches "${MANIFEST}" "${NAME}" "${DEST}" \ + || die "somme de contrôle incorrecte après téléchargement : ${NAME}" +} + +fetch_assets () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local ASSETS="${3}" + local MANIFEST="${4}" + local unit short FULL_NAME + + info "assets de la release ${TAG} pour le profil ${PROFILE}" + + run "${DRY_RUN}" "mkdir -p '${BIN_PATH}'" + run "${DRY_RUN}" "mkdir -p '${UNIT_PATH}'" + run "${DRY_RUN}" "mkdir -p '${LN_PATH}'" + + for unit in $(profile_units "${PROFILE}") + do + fetch_asset "${DRY_RUN}" "${MANIFEST}" "${ASSETS}" "${unit}" "${UNIT_PATH}${unit}" + done + + for short in $(profile_binaries "${PROFILE}") + do + FULL_NAME=$(asset_name "${ASSETS}" "${short}") + [[ -n "${FULL_NAME}" ]] || die "exécutable absent de la release ${TAG} : ${short}" + fetch_asset "${DRY_RUN}" "${MANIFEST}" "${ASSETS}" "${FULL_NAME}" "${BIN_PATH}${FULL_NAME}" + run "${DRY_RUN}" "chmod +x '${BIN_PATH}${FULL_NAME}'" + done +} + +install_units () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local UNITS unit + + UNITS=$(profile_units "${PROFILE}") || die "profil inconnu : ${PROFILE}" + [[ -n "${UNITS}" ]] || { info "units : aucune pour le profil ${PROFILE}"; return 0; } + + info "units du profil ${PROFILE}" + + for unit in ${UNITS} + do + if [[ ${DRY_RUN} -ne ${FLAGS_TRUE} ]] && [[ ! -f "${UNIT_PATH}${unit}" ]] + then + die "unit absente des assets de la release ${TAG} : ${unit}" + fi + run "${DRY_RUN}" "install -m 0644 '${UNIT_PATH}${unit}' '${UNIT_DIR}/${unit}'" + done + + run "${DRY_RUN}" "systemctl daemon-reload" + + for unit in ${UNITS} + do + case "${unit}" in + *@.service) continue ;; + esac + run "${DRY_RUN}" "systemctl enable '${unit}'" + done +} + +switch_binaries () { + local DRY_RUN="${1}" + local PROFILE="${2}" + local ASSETS="${3}" + local BINARIES INSTANCES short FULL_NAME + + BINARIES=$(profile_binaries "${PROFILE}") + [[ -n "${BINARIES}" ]] || { info "bascule : aucun exécutable pour le profil ${PROFILE}"; return 0; } + + info "bascule des binaires" + + INSTANCES=$(active_instances "${PROFILE}") + + stop_services "${DRY_RUN}" "${PROFILE}" "${INSTANCES}" + + for short in ${BINARIES} + do + FULL_NAME=$(asset_name "${ASSETS}" "${short}") + run "${DRY_RUN}" "rm -f '${LN_PATH}${short}'" + run "${DRY_RUN}" "ln -s '${BIN_PATH}${FULL_NAME}' '${LN_PATH}${short}'" + done + + start_services "${DRY_RUN}" "${PROFILE}" "${INSTANCES}" +} + main () { [[ -f ./libs/shflags ]] && . ./libs/shflags || eval "$(curl --silent https://git.g3e.fr/H6N/tools/raw/branch/main/libs/shflags)" - DEFINE_boolean 'dryrun' false 'Enable dry-run mode' 'd' - DEFINE_boolean 'up_script' true 'Upgrade script' 's' - DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g' - DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r' - DEFINE_string 'branch' 'main/' 'Branch name' 'b' - DEFINE_string 'tag' '' 'Tag name' 't' + DEFINE_boolean 'dryrun' false 'Enable dry-run mode' 'd' + DEFINE_boolean 'up_script' true 'Upgrade script' 's' + DEFINE_string 'git_server' 'https://git.g3e.fr/' 'Git Server' 'g' + DEFINE_string 'repo_path' 'syonad/two/' 'Path of repository' 'r' + DEFINE_string 'branch' 'main/' 'Branch name' 'b' + DEFINE_string 'tag' '' 'Tag name' 't' + DEFINE_string 'profile' 'kvm' 'Host profile: kvm' 'p' + DEFINE_boolean 'bootstrap' false 'Prepare the host (packages, kernel, network)' 'i' + DEFINE_boolean 'packages' true 'Install profile packages during --bootstrap' 'k' + DEFINE_boolean 'network' true 'Configure bridges during --bootstrap' 'n' + DEFINE_string 'uplink' 'eno1' 'Physical uplink interface' 'u' + DEFINE_string 'bridge' 'br-000000' 'Main bridge, uplink is enslaved to it' 'B' + DEFINE_string 'pub_bridge' 'br-public' 'Reserved empty bridge' 'P' + DEFINE_integer 'rollback' 120 'Rollback reboot delay in seconds' 'R' + DEFINE_boolean 'verify' true 'Check assets against SHA256SUMS' 'V' FLAGS "$@" || exit $? eval set -- "${FLAGS_ARGV}" - SCRIPT_URL="${FLAGS_git_server}${FLAGS_repo_path}raw/branch/${FLAGS_branch}${SCRIPT_PATH}" - check_latest_script "${SCRIPT_URL}" "${0}" || ( - [[ ${FLAGS_up_script} -eq ${FLAGS_TRUE} ]] && \ - exec_with_dry_run "${FLAGS_dryrun}" "curl --silent '${SCRIPT_URL}' -o '${0}'" - exit 1 - ) + profile_units "${FLAGS_profile}" >/dev/null || die "profil inconnu : ${FLAGS_profile}" - download_binaries "${FLAGS_dryrun}" "${FLAGS_tag}" "${FLAGS_git_server}" "${FLAGS_repo_path}" + if [[ ${FLAGS_up_script} -eq ${FLAGS_TRUE} ]] + then + local SCRIPT_URL="${FLAGS_git_server}${FLAGS_repo_path}raw/branch/${FLAGS_branch}${SCRIPT_PATH}" + if ! check_latest_script "${SCRIPT_URL}" "${0}" + then + run "${FLAGS_dryrun}" "curl --silent '${SCRIPT_URL}' -o '${0}'" + die "script local différent de la branche ${FLAGS_branch} — mis à jour, relancer" + fi + fi + + [[ ${FLAGS_bootstrap} -eq ${FLAGS_TRUE} ]] && \ + run_bootstrap "${FLAGS_dryrun}" "${FLAGS_profile}" "${FLAGS_git_server}" "${FLAGS_repo_path}" "${FLAGS_branch}" + + local TAG BIN_PATH UNIT_PATH LN_PATH ASSETS MANIFEST + + TAG=$(resolve_tag "${FLAGS_tag}" "${FLAGS_git_server}" "${FLAGS_repo_path}") + [[ -n "${TAG}" && "${TAG}" != "null" ]] || die "impossible de déterminer la release à déployer" + + BIN_PATH="/opt/two/${TAG}/bin/" + UNIT_PATH="/opt/two/${TAG}/units/" + LN_PATH="/opt/two/bin/" + + ASSETS=$(release_assets "${FLAGS_git_server}" "${FLAGS_repo_path}" "${TAG}") + [[ -n "${ASSETS}" ]] || die "aucun asset dans la release ${TAG}" + + # Manifeste absent = release antérieure à sa mise en place, ou CI incomplète. + # C'est bloquant : une vérification qui se désactive d'elle-même ne vérifie + # rien. --noverify est la sortie explicite. + MANIFEST="" + if [[ ${FLAGS_verify} -eq ${FLAGS_TRUE} ]] + then + MANIFEST=$(release_manifest "${ASSETS}") \ + || die "${MANIFEST_NAME} absent de la release ${TAG} — --noverify pour déployer sans vérification" + else + warn "vérification des sommes de contrôle désactivée (--noverify)" + fi + + fetch_assets "${FLAGS_dryrun}" "${FLAGS_profile}" "${ASSETS}" "${MANIFEST}" + + install_units "${FLAGS_dryrun}" "${FLAGS_profile}" + switch_binaries "${FLAGS_dryrun}" "${FLAGS_profile}" "${ASSETS}" } [[ "${BASH_SOURCE[0]}" == "${0}" ]] && (main "$@" || exit 1) -[[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1) \ No newline at end of file +[[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1)