diff --git a/.forgejo/workflows/prerelease.yml b/.forgejo/workflows/prerelease.yml index 242b5de..1869374 100644 --- a/.forgejo/workflows/prerelease.yml +++ b/.forgejo/workflows/prerelease.yml @@ -42,69 +42,27 @@ jobs: goarch: ${{ matrix.goarch }} binari: ${{ matrix.binaries }} secrets: inherit - # 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: + upload-scripts: runs-on: docker needs: [set-release-target] strategy: matrix: - 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 + script: + - run-dnsmasq-in-netns.sh steps: - uses: actions/checkout@v3 - name: Move asset run: | mkdir -p "dist" - cp "${{ matrix.path }}" dist/ - - name: Upload asset + cp scripts/${{ matrix.script }} dist/ + - name: Upload script uses: actions/upload-artifact@v3 with: - 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 + name: ${{ matrix.script }}-${{ needs.set-release-target.outputs.release_cible }} + path: dist/${{ matrix.script }} prerelease: runs-on: docker - # 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] + needs: [set-release-target, build] 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 0861353..7e9c461 100644 --- a/api/agent.yaml +++ b/api/agent.yaml @@ -91,7 +91,7 @@ paths: "404": $ref: "#/components/responses/NotFound" "409": - description: VPC not deletable — only running or error states can be deleted, and all its subnets must be deleted first + description: VPC not in a deletable state content: application/json: schema: @@ -146,7 +146,7 @@ paths: schema: $ref: "#/components/schemas/Error" "422": - description: Subnet not found, or not in creating/running state + description: Subnet not found or not in created state content: application/json: schema: @@ -185,12 +185,6 @@ 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" @@ -280,12 +274,6 @@ 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" @@ -326,8 +314,8 @@ components: example: vp-00001 state: type: string - enum: [creating, running, error, deleting, deleted] - example: running + enum: [creating, created, deleting, deleted] + example: created cidr: type: string example: "10.0.0.0/16" @@ -385,8 +373,8 @@ components: example: sn-00001 state: type: string - enum: [creating, running, error, deleting, deleted] - example: running + enum: [creating, created, deleting, deleted] + example: created vpc: type: string example: vpc1 @@ -484,8 +472,8 @@ components: example: vm-00001 state: type: string - enum: [creating, running, error, deleting, deleted] - example: running + enum: [starting, started, stopping, stopped] + example: started metadata_port: type: string example: "80" diff --git a/cmd/agent/main.go b/cmd/agent/main.go index 0425d4e..7083076 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -8,7 +8,6 @@ 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" @@ -32,13 +31,6 @@ 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 f0671f0..36ee9d7 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -68,13 +68,7 @@ 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 { - // 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) - } + w.WriteHeader(http.StatusNotFound) 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 77ab5e4..36e42f7 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") - kv.AddInDB(db, "subnet/sn-exist/state", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "subnet/sn-exist/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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 != "running" { + if result.Name != "sn-1" || result.State != "created" { 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", "running") + kv.AddInDB(db, "subnet/sn-del/state", "created") req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-del", nil) w := httptest.NewRecorder() s.SubnetByNameHandler(w, req) @@ -275,17 +275,6 @@ 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) @@ -298,7 +287,7 @@ func TestDeleteSubnet_NotFound(t *testing.T) { func TestSubnetByName_InvalidMethod(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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 db8bbd5..3af5968 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": "running", + "vm/vm-1/state": "started", "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": "running", + "vm/vm-2/state": "started", "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": "running", + "vm/vm-3/state": "started", "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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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 b5ba517..3516305 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", "running") + kv.AddInDB(db, "vpc/v1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-exist/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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 != "running" { + if result.Name != "vpc-1" || result.State != "created" { 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", "running") + kv.AddInDB(db, "vpc/vpc-del/state", "created") req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-del", nil) w := httptest.NewRecorder() s.VpcByNameHandler(w, req) @@ -186,37 +186,10 @@ 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", "running") - kv.AddInDB(db, "subnet/sn-1/state", "running") + kv.AddInDB(db, "vpc/vpc-busy/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "created") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") req := httptest.NewRequest(http.MethodDelete, "/vpcs/vpc-busy", nil) w := httptest.NewRecorder() @@ -228,7 +201,7 @@ func TestDeleteVpc_BlockedByActiveSubnet(t *testing.T) { func TestVpcByName_InvalidMethod(t *testing.T) { s, db := newTestServer(t) - kv.AddInDB(db, "vpc/vpc-1/state", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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 e4129a2..2897c8a 100644 --- a/internal/dispatcher/agent/dispatcher.go +++ b/internal/dispatcher/agent/dispatcher.go @@ -6,7 +6,6 @@ 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" ) @@ -14,7 +13,6 @@ 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 { @@ -45,10 +43,6 @@ 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 4475a82..b4881d9 100644 --- a/internal/dispatcher/agent/dispatcher_test.go +++ b/internal/dispatcher/agent/dispatcher_test.go @@ -4,10 +4,8 @@ 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" ) @@ -63,66 +61,3 @@ 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 c14f6a3..2cee6fc 100644 --- a/internal/dispatcher/agent/helpers_test.go +++ b/internal/dispatcher/agent/helpers_test.go @@ -25,18 +25,10 @@ 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 6033a40..e488842 100644 --- a/internal/dispatcher/agent/subnet_commands.go +++ b/internal/dispatcher/agent/subnet_commands.go @@ -6,7 +6,6 @@ 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" @@ -23,8 +22,6 @@ 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" @@ -35,18 +32,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 := state.Get(db, "vpc/"+c.VPC) + vpcState, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state") if err != nil { return fmt.Errorf("vpc %q not found", c.VPC) } - if vpcState != state.Creating && vpcState != state.Running { + if vpcState == "deleting" || vpcState == "deleted" { return fmt.Errorf("vpc %q is %s", c.VPC, vpcState) } localIface, ok := cfg.Interfaces[c.IfaceType] if !ok { localIface = cfg.DefaultInterface } - state.Set(db, c.Key(), state.Creating) + kv.AddInDB(db, "subnet/"+c.Name+"/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) @@ -62,19 +59,16 @@ 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 { - vpcState, err := state.Get(db, "vpc/"+c.VPC) + state, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state") if err != nil { return fmt.Errorf("vpc %q not found while waiting", c.VPC) } - if vpcState == state.Running { + if state == "created" { 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 running", c.VPC) + return fmt.Errorf("timed out waiting for vpc %q to be created", c.VPC) case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second): } } @@ -85,28 +79,22 @@ type DeleteSubnetCommand struct { Name string } -func (c DeleteSubnetCommand) Key() string { return "subnet/" + c.Name } - func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - current, err := state.Get(db, c.Key()) - if err != nil { + if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil { return fmt.Errorf("subnet %q not found", c.Name) } - 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) + return kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting") } func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error { if err := subnet.DeleteSubnet(db, c.Name); err != nil { return err } - current, err := state.Get(db, c.Key()) + state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state") if err != nil { return err } - if current == state.Deleted { + if 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 645f620..b6aacee 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") - kv.AddInDB(db, "subnet/sn-exist/state", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "subnet/sn-exist/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-del/state", "created") cmd := DeleteSubnetCommand{Name: "sn-del"} if err := cmd.Prepare(db, nil); err != nil { t.Fatalf("Prepare a échoué : %v", err) @@ -222,32 +222,6 @@ 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 6b50629..bdfeaec 100644 --- a/internal/dispatcher/agent/vm_commands.go +++ b/internal/dispatcher/agent/vm_commands.go @@ -8,7 +8,6 @@ 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" @@ -31,24 +30,22 @@ 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 := state.Get(db, "subnet/"+c.Subnet) + subnetState, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state") if err != nil { return fmt.Errorf("subnet %q not found", c.Subnet) } - if subnetState != state.Creating && subnetState != state.Running { + if subnetState == "deleting" || subnetState == "deleted" { 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) } - state.Set(db, c.Key(), state.Creating) + kv.AddInDB(db, "vm/"+c.Name+"/state", "starting") 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)) @@ -94,19 +91,16 @@ 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 { - subnetState, err := state.Get(db, "subnet/"+c.Subnet) + state, err := kv.GetFromDB(db, "subnet/"+c.Subnet+"/state") if err != nil { return fmt.Errorf("subnet %q not found while waiting", c.Subnet) } - if subnetState == state.Running { + if state == "created" { 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 running", c.Subnet) + return fmt.Errorf("timed out waiting for subnet %q to be created", c.Subnet) case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second): } } @@ -117,28 +111,22 @@ type StopVMCommand struct { Name string } -func (c StopVMCommand) Key() string { return "vm/" + c.Name } - func (c StopVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - current, err := state.Get(db, c.Key()) - if err != nil { + if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err != nil { return fmt.Errorf("vm %q not found", c.Name) } - 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) + return kv.AddInDB(db, "vm/"+c.Name+"/state", "stopping") } func (c StopVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error { if err := vm.StopVM(db, c.Name, cfg); err != nil { return err } - current, err := state.Get(db, c.Key()) + state, err := kv.GetFromDB(db, "vm/"+c.Name+"/state") if err != nil { return err } - if current == state.Deleted { + if state == "stopped" { 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 2747b0f..ba400b9 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") 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", "running") + kv.AddInDB(db, "subnet/sn-1/state", "created") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") cmd := StartVMCommand{ @@ -114,54 +114,9 @@ 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", "running") + kv.AddInDB(db, "vm/vm-exist/state", "started") cmd := StartVMCommand{ Name: "vm-exist", diff --git a/internal/dispatcher/agent/vpc_commands.go b/internal/dispatcher/agent/vpc_commands.go index 9e26819..a129d88 100644 --- a/internal/dispatcher/agent/vpc_commands.go +++ b/internal/dispatcher/agent/vpc_commands.go @@ -7,7 +7,6 @@ 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" @@ -18,8 +17,6 @@ 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) @@ -30,7 +27,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 state.Set(db, c.Key(), state.Creating) + return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating") } func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { @@ -41,16 +38,10 @@ type DeleteVPCCommand struct { Name string } -func (c DeleteVPCCommand) Key() string { return "vpc/" + c.Name } - func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - current, err := state.Get(db, c.Key()) - if err != nil { + if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); 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) @@ -60,12 +51,12 @@ func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error continue } subnetName := strings.Split(key, "/")[1] - s, err := state.Get(db, "subnet/"+subnetName) - if err != nil || (s != state.Deleting && s != state.Deleted) { + state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") + if err != nil || (state != "deleting" && state != "deleted") { return fmt.Errorf("subnet %q must be deleted before deleting vpc %q", subnetName, c.Name) } } - return state.Set(db, c.Key(), state.Deleting) + return kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting") } func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) error { @@ -79,8 +70,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] - s, _ := state.Get(db, "subnet/"+subnetName) - if s == state.Deleting { + state, _ := kv.GetFromDB(db, "subnet/"+subnetName+"/state") + if state == "deleting" { pending = true break } @@ -98,11 +89,11 @@ func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) erro if err := vpc.DeleteVPC(db, c.Name); err != nil { return err } - current, err := state.Get(db, c.Key()) + state, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state") if err != nil { return err } - if current == state.Deleted { + if 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 78b3d98..9633733 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", "running") + kv.AddInDB(db, "vpc/vpc-exist/state", "created") 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", "running") + kv.AddInDB(db, "vpc/vpc-del/state", "created") cmd := DeleteVPCCommand{Name: "vpc-del"} if err := cmd.Prepare(db, nil); err != nil { t.Fatalf("Prepare a échoué : %v", err) @@ -70,45 +70,10 @@ 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", "running") - kv.AddInDB(db, "subnet/sn-1/state", "running") + kv.AddInDB(db, "vpc/vpc-busy/state", "created") + kv.AddInDB(db, "subnet/sn-1/state", "created") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-busy") cmd := DeleteVPCCommand{Name: "vpc-busy"} if err := cmd.Prepare(db, nil); err == nil { @@ -118,7 +83,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", "running") + kv.AddInDB(db, "vpc/vpc-ok/state", "created") kv.AddInDB(db, "subnet/sn-1/state", "deleted") kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-ok") cmd := DeleteVPCCommand{Name: "vpc-ok"} @@ -129,7 +94,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", "running") + kv.AddInDB(db, "vpc/vpc-ok/state", "created") 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 deleted file mode 100644 index 99b5517..0000000 --- a/internal/migration/state.go +++ /dev/null @@ -1,84 +0,0 @@ -// 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 deleted file mode 100644 index 85137a4..0000000 --- a/internal/migration/state_test.go +++ /dev/null @@ -1,175 +0,0 @@ -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 073dd4b..04712ad 100644 --- a/internal/prometheus/agent/collector.go +++ b/internal/prometheus/agent/collector.go @@ -3,13 +3,12 @@ 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 = state.All() +var allStates = []string{"creating", "created", "deleting", "deleted"} // AgentCollector implements prometheus.Collector and exposes agent metrics // by querying the BadgerDB on each scrape. @@ -17,7 +16,6 @@ type AgentCollector struct { db *badger.DB vpcsTotal *prometheus.Desc subnetsTotal *prometheus.Desc - vmsTotal *prometheus.Desc } func NewAgentCollector(db *badger.DB) *AgentCollector { @@ -33,30 +31,23 @@ 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[state.State]float64, len(allStates)) + counts := make(map[string]float64, len(allStates)) for _, s := range allStates { counts[s] = 0 } @@ -64,18 +55,13 @@ 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") { - 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]++ + if strings.HasSuffix(key, "/state") { + counts[val]++ } } } - for _, s := range allStates { - ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[s], string(s)) + for _, state := range allStates { + ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[state], state) } } diff --git a/internal/prometheus/agent/collector_test.go b/internal/prometheus/agent/collector_test.go deleted file mode 100644 index 49272f9..0000000 --- a/internal/prometheus/agent/collector_test.go +++ /dev/null @@ -1,151 +0,0 @@ -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 2177066..9b50add 100644 --- a/internal/qemu/config.go +++ b/internal/qemu/config.go @@ -1,9 +1,5 @@ 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 fa7bcf3..181bb1d 100644 --- a/internal/qemu/start_linux.go +++ b/internal/qemu/start_linux.go @@ -103,16 +103,9 @@ func Start(cfg Config) error { "-daemonize", ) - 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))) + cmd := exec.Command("qemu-system-x86_64", args...) + if err := cmd.Run(); err != nil { + return fmt.Errorf("qemu-system-x86_64: %w", err) } return nil } diff --git a/internal/state/state.go b/internal/state/state.go deleted file mode 100644 index c2fc9f4..0000000 --- a/internal/state/state.go +++ /dev/null @@ -1,68 +0,0 @@ -// 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 deleted file mode 100644 index bcc0cf4..0000000 --- a/internal/state/state_test.go +++ /dev/null @@ -1,168 +0,0 @@ -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 75ef62e..a3faf46 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/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/systemd" "github.com/dgraph-io/badger/v4" ) func CreateSubnet(db *badger.DB, subnetName string) error { - current, err := state.Get(db, "subnet/"+subnetName) + state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") if err != nil { return err } - if current != state.Creating { + if state != "creating" { return nil } @@ -31,7 +31,7 @@ func CreateSubnet(db *badger.DB, subnetName string) error { return err } - return state.Set(db, "subnet/"+subnetName, state.Running) + return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created") } func createSubnet(db *badger.DB, subnetName string, d subnetData) error { diff --git a/internal/subnet/delete.go b/internal/subnet/delete.go index 1175fd5..f745eda 100644 --- a/internal/subnet/delete.go +++ b/internal/subnet/delete.go @@ -7,7 +7,6 @@ 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" @@ -15,11 +14,11 @@ import ( ) func DeleteSubnet(db *badger.DB, subnetName string) error { - current, err := state.Get(db, "subnet/"+subnetName) + state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state") if err != nil { return err } - if current != state.Deleting { + if state != "deleting" { return nil } @@ -45,7 +44,7 @@ func DeleteSubnet(db *badger.DB, subnetName string) error { return fmt.Errorf("unknown subnet mode %q", d.mode) } - return state.Set(db, "subnet/"+subnetName, state.Deleted) + return kv.AddInDB(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 74ed0bb..3e0d4a7 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/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" ) func StartVM(db *badger.DB, name string, cfg *configuration.Config) error { - current, err := state.Get(db, "vm/"+name) + state, err := kv.GetFromDB(db, "vm/"+name+"/state") if err != nil { return err } - if current != state.Creating { + if state != "starting" { 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 state.Set(db, "vm/"+name, state.Running) + return kv.AddInDB(db, "vm/"+name+"/state", "started") } func copyFile(src, dst string) error { diff --git a/internal/vm/delete.go b/internal/vm/delete.go index f1490f0..3d808e6 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/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" ) func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { - current, err := state.Get(db, "vm/"+name) + state, err := kv.GetFromDB(db, "vm/"+name+"/state") if err != nil { return err } - if current != state.Deleting { + if state != "stopping" { return nil } @@ -65,7 +65,7 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { os.Remove(varsPath) } - return state.Set(db, "vm/"+name, state.Deleted) + return kv.AddInDB(db, "vm/"+name+"/state", "stopped") } func waitQMPDead(socketPath string, timeout, poll time.Duration) { diff --git a/internal/vpc/create.go b/internal/vpc/create.go index 52e2567..148f70c 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/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" ) func CreateVPC(db *badger.DB, name string) error { - if current, err := state.Get(db, "vpc/"+name); err != nil { + if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil { return err - } else if current == state.Creating { + } else if 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 } - return state.Set(db, "vpc/"+name, state.Running) + kv.AddInDB(db, "vpc/"+name+"/state", "created") } return nil } diff --git a/internal/vpc/delete.go b/internal/vpc/delete.go index ff4d83f..dbd3a59 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/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" ) func DeleteVPC(db *badger.DB, name string) error { - if current, err := state.Get(db, "vpc/"+name); err != nil { + if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil { return err - } else if current == state.Deleting { + } else if 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 } - return state.Set(db, "vpc/"+name, state.Deleted) + kv.AddInDB(db, "vpc/"+name+"/state", "deleted") } return nil diff --git a/scripts/bootstrap_kvm.sh b/scripts/bootstrap_kvm.sh deleted file mode 100755 index e239b7c..0000000 --- a/scripts/bootstrap_kvm.sh +++ /dev/null @@ -1,169 +0,0 @@ -#!/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 100755 new mode 100644 index 07d7668..7532b14 --- a/scripts/deploy.sh +++ b/scripts/deploy.sh @@ -10,14 +10,6 @@ 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 @@ -26,7 +18,7 @@ exec_with_dry_run () { eval "${2}" 2> /tmp/error || \ { echo -e "failed with following error"; - local output; output=$(cat /tmp/error | sed -e "s/^/ error -> /g"); + output=$(cat /tmp/error | sed -e "s/^/ error -> /g"); echo -e "${output}"; return 1; } @@ -34,358 +26,67 @@ exec_with_dry_run () { return 0 } -run () { - exec_with_dry_run "${1}" "${2}" || die "échec : ${2}" -} - check_latest_script () { - local REMOTE_URL="${1}" - local LOCAL_PATH="${2}" - local REMOTE_SUM LOCAL_SUM + REMOTE_URL="${1}" + LOCAL_PATH="${2}" - REMOTE_SUM=$(curl --silent "${REMOTE_URL}" | sha256sum) - LOCAL_SUM=$(cat ${LOCAL_PATH} | sha256sum) + REMOTE=$(curl --silent "${REMOTE_URL}" | sha256sum) + LOCAL=$(cat ${LOCAL_PATH} | sha256sum) - [[ "${REMOTE_SUM}" == "${LOCAL_SUM}" ]] || return 1 + [[ "${REMOTE}" == "${LOCAL}" ]] || return 1 return 0 } -list_active_units () { - local PATTERN="${1}" - systemctl list-units --state=active --no-legend --plain "${PATTERN}" 2>/dev/null \ - | awk '{print $1}' || true -} +download_binaries () { + DRY_RUN="${1}" + TAG="${2}" + GIT_SERVER="${3}" + REPO_PATH="${4}" -# 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}") + #'.[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 do - case "${unit}" in - *@.service) continue ;; - *) echo "${unit}" ;; - esac + 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}'" 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_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' + 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' FLAGS "$@" || exit $? eval set -- "${FLAGS_ARGV}" - profile_units "${FLAGS_profile}" >/dev/null || die "profil inconnu : ${FLAGS_profile}" + 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 + ) - 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}" + download_binaries "${FLAGS_dryrun}" "${FLAGS_tag}" "${FLAGS_git_server}" "${FLAGS_repo_path}" } [[ "${BASH_SOURCE[0]}" == "${0}" ]] && (main "$@" || exit 1) -[[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1) +[[ "${BASH_SOURCE[0]}" == "" ]] && (main "$@" || exit 1) \ No newline at end of file