diff --git a/learn/projects/sw-block/phases/phase-5-dev-log.md b/learn/projects/sw-block/phases/phase-5-dev-log.md new file mode 100644 index 000000000..79f0c0c79 --- /dev/null +++ b/learn/projects/sw-block/phases/phase-5-dev-log.md @@ -0,0 +1,36 @@ +# Phase 5 Dev Log + +Append-only communication between agents. Newest entries at bottom. +Each entry: `[date] [role] message` + +Roles: `DEV`, `REVIEWER`, `TESTER`, `ARCHITECT` + +--- + +[2026-03-03] [DEV] CP5-1 ALUA + multipath complete. Added ALUA provider + REPORT TPG (implicit ALUA), VPD 0x83 +NAA+TPG+RTP descriptors, TPGS=01 in INQUIRY, standby write fencing, and -tpg-id flag. Added UUID to VolumeInfo for +shared NAA. Added multipath config and setup script. 4 multipath integration tests added. 10 ALUA unit tests added +(SCSI tests total 53). Reviewer fixes applied: RoleNone maps to Active/Optimized to avoid single-node regression; +REPORT TPG advertises T_SUP when state is Transitioning; TPG ID validation; non-ASCII log fix. Added 2 tests: +alua_role_none_allows_writes and alua_report_tpg_transitioning. All unit tests pass, Linux cross-compile verified. + +[2026-03-03] [TESTER] CP5-1 adversarial suite: 16 tests added/validated (state boundaries, VPD 0x83, REPORT TPG, +concurrency, INQUIRY invariants). All 16 PASS. No regressions in engine + iSCSI tests. + +[2026-03-03] [DEV] CP5-2 CoW snapshots completed. Fixes applied from review: DeleteSnapshot pauses flusher before +closing delta; RestoreSnapshot checks PauseAndFlush error + defers Resume; CreateSnapshot holds snapMu across check/insert; +Delete/Restore use beginOp/endOp; lock order documented (flushMu -> snapMu); non-ASCII punctuation removed; persistSuperblock +now returns error and callers propagate. All tests passing (known pre-existing flaky +rebuild_full_extent_midcopy_writes under full-suite load). + +[2026-03-03] [TESTER] CP5-2 QA adversarial suite: 22 tests in 5 groups (races, role rejection, edge cases, lifecycle, +restore correctness) all PASS. Confirms fixes for delete_during_flush_cow, concurrent_create_same_id, and restore path +nextLSN reset. + +[2026-03-03] [DEV] CP5-3 implementation complete. CHAP: ValidateCHAPConfig with ErrCHAPSecretEmpty and CLI guard +requires -chap-secret when -chap-user is set. Login SecurityNeg echoes AuthMethod=CHAP on second PDU after verify; test +assertion added. Metrics adapter docs clarify counters count attempts; /metrics inherits admin auth noted in header +comment. All CP5-3 tests pass; only pre-existing flaky rebuild_catchup_concurrent_writes observed under full suite. + +[2026-03-03] [TESTER] CP5-3 QA adversarial: 28 tests added (16 CHAP + 12 resize) all PASS. No new bugs. Full +regression clean except pre-existing flaky rebuild_catchup_concurrent_writes. diff --git a/learn/projects/sw-block/phases/phase-5-progress.md b/learn/projects/sw-block/phases/phase-5-progress.md new file mode 100644 index 000000000..bcae8f8a2 --- /dev/null +++ b/learn/projects/sw-block/phases/phase-5-progress.md @@ -0,0 +1,32 @@ +# Phase 5 Progress + +## Status +- CP5-1 ALUA + multipath complete. CP5-2 CoW snapshots complete. CP5-3 complete. + +## Completed +- CP5-1: ALUA implicit support, REPORT TARGET PORT GROUPS, VPD 0x83 descriptors, write fencing on standby. +- CP5-1: Multipath config + setup script, 4 multipath integration tests. +- CP5-1: Reviewer fixes (RoleNone write regression, T_SUP flag, TPG ID validation, ASCII log). +- CP5-1: 10 ALUA unit tests + 16 adversarial tests (all PASS). +- CP5-2: CoW snapshots implemented with flusher-based CoW, delta files, and recovery. +- CP5-2: Review fixes applied (PauseAndFlush safety, snapMu race fix, beginOp/endOp, lock order doc, error propagation). +- CP5-2: 10 unit tests + 22 adversarial tests (all PASS). +- CP5-3: CHAP auth, online resize, Prometheus metrics, admin endpoints. +- CP5-3: Review fixes applied (empty secret validation, AuthMethod echo, docs). +- CP5-3: 12 dev tests + 28 QA adversarial tests (all PASS). + +## In Progress +- CP5-4: Failure injection + Layer-5 validation (not started). + +## Blockers +- None. + +## Next Steps +- Decide CP5-2 scope (CSI driver vs CHAP/metrics/admin CLI). + +## Notes +- SCSI test count: 53 (12 ALUA). Integration multipath tests require multipath-tools + sg3_utils. +- Known flaky: rebuild_full_extent_midcopy_writes under full-suite CPU contention (pre-existing). +- Known flaky: rebuild_catchup_concurrent_writes (WAL_RECYCLED timing, pre-existing). + + diff --git a/weed/storage/blockvol/blockvol.go b/weed/storage/blockvol/blockvol.go index ebb607349..828c17e9e 100644 --- a/weed/storage/blockvol/blockvol.go +++ b/weed/storage/blockvol/blockvol.go @@ -541,6 +541,22 @@ type VolumeInfo struct { Healthy bool } +// WALUsedFraction returns the fraction of WAL space currently in use (0.0 to 1.0). +func (v *BlockVol) WALUsedFraction() float64 { + if v.wal == nil { + return 0 + } + return v.wal.UsedFraction() +} + +// DirtyMapLen returns the number of entries in the dirty map. +func (v *BlockVol) DirtyMapLen() int { + if v.dirtyMap == nil { + return 0 + } + return v.dirtyMap.Len() +} + // SyncCache ensures all previously written WAL entries are durable on disk. // It submits a sync request to the group committer, which batches fsyncs. func (v *BlockVol) SyncCache() error { @@ -862,6 +878,58 @@ func (v *BlockVol) ListSnapshots() []SnapshotInfo { return infos } +var ( + ErrShrinkNotSupported = errors.New("blockvol: shrink not supported") + ErrSnapshotsPreventResize = errors.New("blockvol: cannot resize with active snapshots") +) + +// Expand grows the volume to newSize bytes. newSize must be larger than +// the current size and aligned to BlockSize. Fails if snapshots are active. +func (v *BlockVol) Expand(newSize uint64) error { + if err := v.beginOp(); err != nil { + return err + } + defer v.endOp() + if err := v.writeGate(); err != nil { + return err + } + + if newSize <= v.super.VolumeSize { + if newSize == v.super.VolumeSize { + return nil // no-op + } + return ErrShrinkNotSupported + } + if newSize%uint64(v.super.BlockSize) != 0 { + return ErrAlignment + } + + // Hold snapMu across entire operation to prevent concurrent CreateSnapshot. + v.snapMu.RLock() + defer v.snapMu.RUnlock() + if len(v.snapshots) > 0 { + return ErrSnapshotsPreventResize + } + + // Pause flusher (no concurrent extent writes during file extension). + if err := v.flusher.PauseAndFlush(); err != nil { + v.flusher.Resume() + return fmt.Errorf("blockvol: expand flush: %w", err) + } + defer v.flusher.Resume() + + // Extend file. + extentStart := v.super.WALOffset + v.super.WALSize + newFileSize := int64(extentStart + newSize) + if err := v.fd.Truncate(newFileSize); err != nil { + return fmt.Errorf("blockvol: expand truncate: %w", err) + } + + // Update superblock. + v.super.VolumeSize = newSize + return v.persistSuperblock() +} + // persistSuperblock writes the superblock to disk and fsyncs. func (v *BlockVol) persistSuperblock() error { if _, err := v.fd.Seek(0, 0); err != nil { diff --git a/weed/storage/blockvol/iscsi/auth.go b/weed/storage/blockvol/iscsi/auth.go new file mode 100644 index 000000000..29a784b71 --- /dev/null +++ b/weed/storage/blockvol/iscsi/auth.go @@ -0,0 +1,109 @@ +// auth.go implements CHAP authentication for iSCSI (RFC 7143 S12.1). +// Only unidirectional (target authenticates initiator) with MD5 (algorithm 5). +package iscsi + +import ( + "crypto/md5" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "strings" +) + +var ErrCHAPSecretEmpty = errors.New("iscsi: CHAP enabled but secret is empty") + +// CHAPConfig holds CHAP authentication settings. +type CHAPConfig struct { + Enabled bool + Username string // expected initiator username (empty = accept any) + Secret string // shared secret +} + +type chapState int + +const ( + chapIdle chapState = iota + chapChallengeSent + chapDone +) + +// CHAPAuthenticator drives the target side of a CHAP exchange. +type CHAPAuthenticator struct { + config CHAPConfig + state chapState + id uint8 // challenge identifier (0-255) + challenge []byte // 16 random bytes +} + +// ValidateCHAPConfig checks that a CHAPConfig is well-formed. +// Must be called at startup before passing the config to NewLoginNegotiator. +func ValidateCHAPConfig(c CHAPConfig) error { + if c.Enabled && c.Secret == "" { + return ErrCHAPSecretEmpty + } + return nil +} + +// NewCHAPAuthenticator creates a CHAP authenticator for one login session. +// The config must have been validated with ValidateCHAPConfig at startup. +func NewCHAPAuthenticator(config CHAPConfig) *CHAPAuthenticator { + return &CHAPAuthenticator{config: config, state: chapIdle} +} + +// IsEnabled returns whether CHAP authentication is required. +func (a *CHAPAuthenticator) IsEnabled() bool { + return a.config.Enabled +} + +// GenerateChallenge produces the CHAP_A, CHAP_I, CHAP_C parameters for the +// first security negotiation response. Must be called exactly once. +func (a *CHAPAuthenticator) GenerateChallenge() (map[string]string, error) { + // Generate random id byte. + var idBuf [1]byte + if _, err := rand.Read(idBuf[:]); err != nil { + return nil, fmt.Errorf("chap: generate id: %w", err) + } + a.id = idBuf[0] + + // Generate 16-byte random challenge. + a.challenge = make([]byte, 16) + if _, err := rand.Read(a.challenge); err != nil { + return nil, fmt.Errorf("chap: generate challenge: %w", err) + } + + a.state = chapChallengeSent + + return map[string]string{ + "CHAP_A": "5", // MD5 + "CHAP_I": fmt.Sprintf("%d", a.id), // decimal + "CHAP_C": "0x" + hex.EncodeToString(a.challenge), // hex with 0x prefix + }, nil +} + +// Verify checks the initiator's CHAP_N (username) and CHAP_R (response). +// Returns true if authentication succeeds. +func (a *CHAPAuthenticator) Verify(chapN, chapR string) bool { + if a.state != chapChallengeSent { + return false + } + a.state = chapDone + + // Check username if configured. + if a.config.Username != "" && chapN != a.config.Username { + return false + } + + // Compute expected response: MD5(id_byte || secret_bytes || challenge_bytes). + h := md5.New() + h.Write([]byte{a.id}) + h.Write([]byte(a.config.Secret)) + h.Write(a.challenge) + expected := hex.EncodeToString(h.Sum(nil)) + + // Normalize initiator response: strip "0x" prefix if present. + got := strings.TrimPrefix(chapR, "0x") + got = strings.TrimPrefix(got, "0X") + + return strings.EqualFold(expected, got) +} diff --git a/weed/storage/blockvol/iscsi/auth_test.go b/weed/storage/blockvol/iscsi/auth_test.go new file mode 100644 index 000000000..9c8c8741a --- /dev/null +++ b/weed/storage/blockvol/iscsi/auth_test.go @@ -0,0 +1,183 @@ +package iscsi + +import ( + "crypto/md5" + "encoding/hex" + "strconv" + "strings" + "testing" +) + +// TestCHAP_LoginSuccess verifies that a correct CHAP username/password +// completes login through the full SecurityNeg -> LoginOp -> FFP flow. +func TestCHAP_LoginSuccess(t *testing.T) { + config := DefaultTargetConfig() + config.TargetName = "iqn.2024.com.seaweedfs:vol1" + config.CHAPConfig = CHAPConfig{ + Enabled: true, + Username: "testuser", + Secret: "s3cret", + } + + ln := NewLoginNegotiator(config) + resolver := newResolver(config.TargetName) + + // PDU 1: Initiator sends SecurityNeg with AuthMethod=CHAP + p1 := NewParams() + p1.Set("InitiatorName", "iqn.2024.com.test:initiator1") + p1.Set("TargetName", config.TargetName) + p1.Set("AuthMethod", "CHAP") + req1 := makeLoginReq(StageSecurityNeg, StageLoginOp, true, p1) + + resp1 := ln.HandleLoginPDU(req1, resolver) + if resp1.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success, got class=%d detail=%d", + resp1.LoginStatusClass(), resp1.LoginStatusDetail()) + } + if resp1.LoginTransit() { + t.Fatal("expected T=0 (no transit) in challenge response") + } + + // Parse challenge params from response + rp1, err := ParseParams(resp1.DataSegment) + if err != nil { + t.Fatalf("parse resp1 params: %v", err) + } + chapA, _ := rp1.Get("CHAP_A") + if chapA != "5" { + t.Fatalf("expected CHAP_A=5, got %q", chapA) + } + chapIStr, _ := rp1.Get("CHAP_I") + chapCStr, _ := rp1.Get("CHAP_C") + if chapIStr == "" || chapCStr == "" { + t.Fatalf("missing CHAP_I or CHAP_C in response") + } + + // Compute CHAP response + chapID, _ := strconv.Atoi(chapIStr) + challenge, _ := hex.DecodeString(strings.TrimPrefix(chapCStr, "0x")) + h := md5.New() + h.Write([]byte{byte(chapID)}) + h.Write([]byte("s3cret")) + h.Write(challenge) + chapR := "0x" + hex.EncodeToString(h.Sum(nil)) + + // PDU 2: Initiator sends CHAP_N + CHAP_R + p2 := NewParams() + p2.Set("CHAP_N", "testuser") + p2.Set("CHAP_R", chapR) + req2 := makeLoginReq(StageSecurityNeg, StageLoginOp, true, p2) + + resp2 := ln.HandleLoginPDU(req2, resolver) + if resp2.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success after CHAP, got class=%d detail=%d", + resp2.LoginStatusClass(), resp2.LoginStatusDetail()) + } + if !resp2.LoginTransit() { + t.Fatal("expected T=1 (transit) after successful CHAP") + } + // Verify AuthMethod=CHAP echoed in the second response. + rp2, err := ParseParams(resp2.DataSegment) + if err != nil { + t.Fatalf("parse resp2 params: %v", err) + } + if am, ok := rp2.Get("AuthMethod"); !ok || am != "CHAP" { + t.Fatalf("expected AuthMethod=CHAP in second response, got %q (ok=%v)", am, ok) + } + + // PDU 3: LoginOp -> FullFeature + p3 := NewParams() + p3.Set("MaxRecvDataSegmentLength", "65536") + req3 := makeLoginReq(StageLoginOp, StageFullFeature, true, p3) + + resp3 := ln.HandleLoginPDU(req3, resolver) + if resp3.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success at FFP, got class=%d detail=%d", + resp3.LoginStatusClass(), resp3.LoginStatusDetail()) + } + if !ln.Done() { + t.Fatal("expected login Done after FFP transition") + } +} + +// TestCHAP_LoginWrongPassword verifies that an incorrect CHAP response +// is rejected with AuthFailure. +func TestCHAP_LoginWrongPassword(t *testing.T) { + config := DefaultTargetConfig() + config.TargetName = "iqn.2024.com.seaweedfs:vol1" + config.CHAPConfig = CHAPConfig{ + Enabled: true, + Username: "testuser", + Secret: "s3cret", + } + + ln := NewLoginNegotiator(config) + resolver := newResolver(config.TargetName) + + // PDU 1: SecurityNeg with AuthMethod=CHAP + p1 := NewParams() + p1.Set("InitiatorName", "iqn.2024.com.test:initiator1") + p1.Set("TargetName", config.TargetName) + p1.Set("AuthMethod", "CHAP") + req1 := makeLoginReq(StageSecurityNeg, StageLoginOp, true, p1) + + resp1 := ln.HandleLoginPDU(req1, resolver) + if resp1.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success for challenge, got class=%d detail=%d", + resp1.LoginStatusClass(), resp1.LoginStatusDetail()) + } + + // PDU 2: Wrong password + p2 := NewParams() + p2.Set("CHAP_N", "testuser") + p2.Set("CHAP_R", "0xdeadbeefdeadbeefdeadbeefdeadbeef") // wrong + req2 := makeLoginReq(StageSecurityNeg, StageLoginOp, true, p2) + + resp2 := ln.HandleLoginPDU(req2, resolver) + if resp2.LoginStatusClass() != LoginStatusInitiatorErr || + resp2.LoginStatusDetail() != LoginDetailAuthFailure { + t.Fatalf("expected AuthFailure, got class=%d detail=%d", + resp2.LoginStatusClass(), resp2.LoginStatusDetail()) + } +} + +// TestCHAP_DisabledAllowsLogin verifies that with CHAP disabled, +// AuthMethod=None works as before. +func TestCHAP_DisabledAllowsLogin(t *testing.T) { + config := DefaultTargetConfig() + config.TargetName = "iqn.2024.com.seaweedfs:vol1" + // CHAPConfig.Enabled defaults to false + + ln := NewLoginNegotiator(config) + resolver := newResolver(config.TargetName) + + // Single SecurityNeg PDU with transit to LoginOp + p := NewParams() + p.Set("InitiatorName", "iqn.2024.com.test:initiator1") + p.Set("TargetName", config.TargetName) + p.Set("AuthMethod", "None") + req := makeLoginReq(StageSecurityNeg, StageLoginOp, true, p) + + resp := ln.HandleLoginPDU(req, resolver) + if resp.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success, got class=%d detail=%d", + resp.LoginStatusClass(), resp.LoginStatusDetail()) + } + if !resp.LoginTransit() { + t.Fatal("expected transit with CHAP disabled") + } + + // LoginOp -> FFP + p2 := NewParams() + p2.Set("MaxRecvDataSegmentLength", "65536") + req2 := makeLoginReq(StageLoginOp, StageFullFeature, true, p2) + + resp2 := ln.HandleLoginPDU(req2, resolver) + if resp2.LoginStatusClass() != LoginStatusSuccess { + t.Fatalf("expected success at FFP, got class=%d detail=%d", + resp2.LoginStatusClass(), resp2.LoginStatusDetail()) + } + if !ln.Done() { + t.Fatal("expected login Done") + } +} diff --git a/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin.go b/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin.go index 39aac99fe..a36fd4887 100644 --- a/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin.go +++ b/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin.go @@ -4,6 +4,9 @@ // GET /status -- return JSON status // POST /replica -- set WAL shipping target {data_addr, ctrl_addr} // POST /rebuild -- start/stop rebuild server {action, listen_addr} +// POST /snapshot -- create/delete/restore/list snapshots +// POST /resize -- expand volume {new_size_bytes} +// GET /metrics -- Prometheus metrics (requires X-Admin-Token if auth enabled) package main import ( @@ -15,14 +18,17 @@ import ( _ "net/http/pprof" // registers /debug/pprof/* handlers on DefaultServeMux "time" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" "github.com/seaweedfs/seaweedfs/weed/storage/blockvol" ) // adminServer provides HTTP admin control of the BlockVol. type adminServer struct { - vol *blockvol.BlockVol - token string // optional auth token; empty = no auth - logger *log.Logger + vol *blockvol.BlockVol + token string // optional auth token; empty = no auth + logger *log.Logger + metricsRegistry *prometheus.Registry // dedicated Prometheus registry (nil = no /metrics) } // assignRequest is the JSON body for POST /assign. @@ -44,6 +50,17 @@ type rebuildRequest struct { ListenAddr string `json:"listen_addr"` } +// snapshotRequest is the JSON body for POST /snapshot. +type snapshotRequest struct { + Action string `json:"action"` // "create", "delete", "restore", "list" + ID uint32 `json:"id"` +} + +// resizeRequest is the JSON body for POST /resize. +type resizeRequest struct { + NewSizeBytes uint64 `json:"new_size_bytes"` +} + // statusResponse is the JSON body for GET /status. type statusResponse struct { Path string `json:"path"` @@ -53,6 +70,7 @@ type statusResponse struct { CheckpointLSN uint64 `json:"checkpoint_lsn"` HasLease bool `json:"has_lease"` Healthy bool `json:"healthy"` + VolumeSize uint64 `json:"volume_size"` } const maxValidRole = uint32(blockvol.RoleDraining) @@ -78,6 +96,12 @@ func (a *adminServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { a.handleReplica(w, r) case "/rebuild": a.handleRebuild(w, r) + case "/snapshot": + a.handleSnapshot(w, r) + case "/resize": + a.handleResize(w, r) + case "/metrics": + a.handleMetrics(w, r) default: http.NotFound(w, r) } @@ -130,6 +154,7 @@ func (a *adminServer) handleStatus(w http.ResponseWriter, r *http.Request) { CheckpointLSN: st.CheckpointLSN, HasLease: st.HasLease, Healthy: info.Healthy, + VolumeSize: info.VolumeSize, } w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(resp) @@ -188,6 +213,98 @@ func (a *adminServer) handleRebuild(w http.ResponseWriter, r *http.Request) { w.Write([]byte(`{"ok":true}`)) } +func (a *adminServer) handleSnapshot(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed) + return + } + var req snapshotRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + jsonError(w, "bad json: "+err.Error(), http.StatusBadRequest) + return + } + + switch req.Action { + case "create": + if err := a.vol.CreateSnapshot(req.ID); err != nil { + jsonError(w, err.Error(), http.StatusConflict) + return + } + a.logger.Printf("admin: created snapshot %d", req.ID) + case "delete": + if err := a.vol.DeleteSnapshot(req.ID); err != nil { + jsonError(w, err.Error(), http.StatusNotFound) + return + } + a.logger.Printf("admin: deleted snapshot %d", req.ID) + case "restore": + if err := a.vol.RestoreSnapshot(req.ID); err != nil { + jsonError(w, err.Error(), http.StatusConflict) + return + } + a.logger.Printf("admin: restored snapshot %d", req.ID) + case "list": + snaps := a.vol.ListSnapshots() + type snapInfo struct { + ID uint32 `json:"id"` + BaseLSN uint64 `json:"base_lsn"` + CreatedAt string `json:"created_at"` + CoWBlocks uint64 `json:"cow_blocks"` + } + result := make([]snapInfo, len(snaps)) + for i, s := range snaps { + result[i] = snapInfo{ + ID: s.ID, + BaseLSN: s.BaseLSN, + CreatedAt: s.CreatedAt.Format(time.RFC3339), + CoWBlocks: s.CoWBlocks, + } + } + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]interface{}{"snapshots": result}) + return + default: + jsonError(w, "action must be 'create', 'delete', 'restore', or 'list'", http.StatusBadRequest) + return + } + w.Header().Set("Content-Type", "application/json") + w.Write([]byte(`{"ok":true}`)) +} + +func (a *adminServer) handleResize(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed) + return + } + var req resizeRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + jsonError(w, "bad json: "+err.Error(), http.StatusBadRequest) + return + } + if req.NewSizeBytes == 0 { + jsonError(w, "new_size_bytes is required", http.StatusBadRequest) + return + } + if err := a.vol.Expand(req.NewSizeBytes); err != nil { + jsonError(w, err.Error(), http.StatusConflict) + return + } + a.logger.Printf("admin: resized to %d bytes", req.NewSizeBytes) + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]interface{}{ + "ok": true, + "volume_size": a.vol.Info().VolumeSize, + }) +} + +func (a *adminServer) handleMetrics(w http.ResponseWriter, r *http.Request) { + if a.metricsRegistry == nil { + http.Error(w, `{"error":"metrics not configured"}`, http.StatusNotFound) + return + } + promhttp.HandlerFor(a.metricsRegistry, promhttp.HandlerOpts{}).ServeHTTP(w, r) +} + // startAdminServer starts the HTTP admin server in a background goroutine. // Returns the listener so tests can determine the actual bound port. // Includes /debug/pprof/* endpoints for profiling. @@ -201,6 +318,9 @@ func startAdminServer(addr string, srv *adminServer) (net.Listener, error) { mux.Handle("/status", srv) mux.Handle("/replica", srv) mux.Handle("/rebuild", srv) + mux.Handle("/snapshot", srv) + mux.Handle("/resize", srv) + mux.Handle("/metrics", srv) // pprof handlers registered on DefaultServeMux by net/http/pprof import. mux.Handle("/debug/pprof/", http.DefaultServeMux) go func() { diff --git a/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin_test.go b/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin_test.go new file mode 100644 index 000000000..b93e883d4 --- /dev/null +++ b/weed/storage/blockvol/iscsi/cmd/iscsi-target/admin_test.go @@ -0,0 +1,151 @@ +package main + +import ( + "bytes" + "encoding/json" + "log" + "net/http" + "os" + "path/filepath" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/storage/blockvol" +) + +// TestAdmin_SnapshotCreateListDelete creates a snapshot via the API, +// lists it, deletes it, and verifies the list is empty. +func TestAdmin_SnapshotCreateListDelete(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "snap-admin.blk") + + vol, err := blockvol.CreateBlockVol(path, blockvol.CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + logger := log.New(os.Stderr, "[test] ", 0) + adm := newAdminServer(vol, "", logger) + ln, err := startAdminServer("127.0.0.1:0", adm) + if err != nil { + t.Fatalf("start admin: %v", err) + } + defer ln.Close() + base := "http://" + ln.Addr().String() + + // Create snapshot 1. + body, _ := json.Marshal(snapshotRequest{Action: "create", ID: 1}) + resp, err := http.Post(base+"/snapshot", "application/json", bytes.NewReader(body)) + if err != nil { + t.Fatalf("POST /snapshot create: %v", err) + } + resp.Body.Close() + if resp.StatusCode != 200 { + t.Fatalf("create snapshot: expected 200, got %d", resp.StatusCode) + } + + // List snapshots. + body2, _ := json.Marshal(snapshotRequest{Action: "list"}) + resp2, err := http.Post(base+"/snapshot", "application/json", bytes.NewReader(body2)) + if err != nil { + t.Fatalf("POST /snapshot list: %v", err) + } + var listResp struct { + Snapshots []struct { + ID uint32 `json:"id"` + } `json:"snapshots"` + } + json.NewDecoder(resp2.Body).Decode(&listResp) + resp2.Body.Close() + if len(listResp.Snapshots) != 1 || listResp.Snapshots[0].ID != 1 { + t.Fatalf("expected 1 snapshot with ID=1, got %+v", listResp.Snapshots) + } + + // Delete snapshot 1. + body3, _ := json.Marshal(snapshotRequest{Action: "delete", ID: 1}) + resp3, err := http.Post(base+"/snapshot", "application/json", bytes.NewReader(body3)) + if err != nil { + t.Fatalf("POST /snapshot delete: %v", err) + } + resp3.Body.Close() + if resp3.StatusCode != 200 { + t.Fatalf("delete snapshot: expected 200, got %d", resp3.StatusCode) + } + + // Verify empty. + body4, _ := json.Marshal(snapshotRequest{Action: "list"}) + resp4, err := http.Post(base+"/snapshot", "application/json", bytes.NewReader(body4)) + if err != nil { + t.Fatalf("POST /snapshot list after delete: %v", err) + } + var listResp2 struct { + Snapshots []struct{} `json:"snapshots"` + } + json.NewDecoder(resp4.Body).Decode(&listResp2) + resp4.Body.Close() + if len(listResp2.Snapshots) != 0 { + t.Fatalf("expected 0 snapshots after delete, got %d", len(listResp2.Snapshots)) + } +} + +// TestAdmin_ResizeExpand resizes a volume via the API and verifies the +// new size in /status. +func TestAdmin_ResizeExpand(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "resize-admin.blk") + + vol, err := blockvol.CreateBlockVol(path, blockvol.CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + logger := log.New(os.Stderr, "[test] ", 0) + adm := newAdminServer(vol, "", logger) + ln, err := startAdminServer("127.0.0.1:0", adm) + if err != nil { + t.Fatalf("start admin: %v", err) + } + defer ln.Close() + base := "http://" + ln.Addr().String() + + // Resize to 2MB. + newSize := uint64(2 * 1024 * 1024) + body, _ := json.Marshal(resizeRequest{NewSizeBytes: newSize}) + resp, err := http.Post(base+"/resize", "application/json", bytes.NewReader(body)) + if err != nil { + t.Fatalf("POST /resize: %v", err) + } + var resizeResp struct { + OK bool `json:"ok"` + VolumeSize uint64 `json:"volume_size"` + } + json.NewDecoder(resp.Body).Decode(&resizeResp) + resp.Body.Close() + if resp.StatusCode != 200 { + t.Fatalf("resize: expected 200, got %d", resp.StatusCode) + } + if !resizeResp.OK || resizeResp.VolumeSize != newSize { + t.Fatalf("resize response: %+v", resizeResp) + } + + // Verify size in /status. + resp2, err := http.Get(base + "/status") + if err != nil { + t.Fatalf("GET /status: %v", err) + } + var status statusResponse + json.NewDecoder(resp2.Body).Decode(&status) + resp2.Body.Close() + if status.VolumeSize != newSize { + t.Fatalf("status VolumeSize: expected %d, got %d", newSize, status.VolumeSize) + } +} diff --git a/weed/storage/blockvol/iscsi/cmd/iscsi-target/main.go b/weed/storage/blockvol/iscsi/cmd/iscsi-target/main.go index 9a382de6e..801299505 100644 --- a/weed/storage/blockvol/iscsi/cmd/iscsi-target/main.go +++ b/weed/storage/blockvol/iscsi/cmd/iscsi-target/main.go @@ -17,6 +17,7 @@ import ( "syscall" "time" + "github.com/prometheus/client_golang/prometheus" "github.com/seaweedfs/seaweedfs/weed/storage/blockvol" "github.com/seaweedfs/seaweedfs/weed/storage/blockvol/iscsi" ) @@ -34,6 +35,8 @@ func main() { replicaData := flag.String("replica-data", "", "replica receiver data listen address (e.g. :9001; empty = disabled)") replicaCtrl := flag.String("replica-ctrl", "", "replica receiver ctrl listen address (e.g. :9002; empty = disabled)") rebuildListen := flag.String("rebuild-listen", "", "rebuild server listen address (e.g. :9003; empty = disabled)") + chapUser := flag.String("chap-user", "", "CHAP username (empty = CHAP disabled)") + chapSecret := flag.String("chap-secret", "", "CHAP shared secret") flag.Parse() if *volPath == "" { @@ -44,6 +47,9 @@ func main() { if *tpgID < 1 || *tpgID > 65535 { log.Fatalf("invalid -tpg-id %d: must be 1-65535", *tpgID) } + if *chapUser != "" && *chapSecret == "" { + log.Fatalf("-chap-secret is required when -chap-user is set") + } logger := log.New(os.Stdout, "[iscsi] ", log.LstdFlags) @@ -102,9 +108,18 @@ func main() { logger.Printf("rebuild server: %s", *rebuildListen) } + // Create Prometheus registry and metrics adapter. + promReg := prometheus.NewRegistry() + instrumented := &instrumentedAdapter{ + inner: &blockVolAdapter{vol: vol, tpgID: uint16(*tpgID)}, + logger: logger, + } + adapter := newMetricsAdapter(instrumented, vol, promReg) + // Start admin HTTP server if configured if *adminAddr != "" { adm := newAdminServer(vol, *adminToken, logger) + adm.metricsRegistry = promReg ln, err := startAdminServer(*adminAddr, adm) if err != nil { log.Fatalf("start admin server: %v", err) @@ -113,16 +128,18 @@ func main() { logger.Printf("admin server: %s", ln.Addr()) } - // Create adapter with ALUA support and latency instrumentation - adapter := &instrumentedAdapter{ - inner: &blockVolAdapter{vol: vol, tpgID: uint16(*tpgID)}, - logger: logger, - } - // Create target server config := iscsi.DefaultTargetConfig() config.TargetName = *iqn config.TargetAlias = "SeaweedFS BlockVol" + if *chapUser != "" && *chapSecret != "" { + config.CHAPConfig = iscsi.CHAPConfig{ + Enabled: true, + Username: *chapUser, + Secret: *chapSecret, + } + logger.Printf("CHAP authentication enabled for user %q", *chapUser) + } if *portal != "" { // Parse portal group tag from "addr:port,tpgt" format if idx := strings.LastIndex(*portal, ","); idx >= 0 { @@ -138,7 +155,7 @@ func main() { ts.AddVolume(*iqn, adapter) // Start periodic performance stats logging (every 5 seconds). - adapter.StartStatsLogger(5 * time.Second) + instrumented.StartStatsLogger(5 * time.Second) // Graceful shutdown on signal sigCh := make(chan os.Signal, 1) diff --git a/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics.go b/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics.go new file mode 100644 index 000000000..00c41e858 --- /dev/null +++ b/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics.go @@ -0,0 +1,194 @@ +package main + +import ( + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/seaweedfs/seaweedfs/weed/storage/blockvol" + "github.com/seaweedfs/seaweedfs/weed/storage/blockvol/iscsi" +) + +// metricsAdapter wraps a BlockDevice and feeds Prometheus counters/histograms. +// It sits in the adapter stack: metricsAdapter -> instrumentedAdapter -> blockVolAdapter. +// Counters count all attempts (including errors) per Prometheus conventions. +type metricsAdapter struct { + inner iscsi.BlockDevice + + writeOps prometheus.Counter + readOps prometheus.Counter + trimOps prometheus.Counter + syncOps prometheus.Counter + writeBytes prometheus.Counter + readBytes prometheus.Counter + writeLatency prometheus.Observer + readLatency prometheus.Observer + syncLatency prometheus.Observer +} + +// gaugeSource provides gauge data from the BlockVol engine. +type gaugeSource struct { + vol *blockvol.BlockVol +} + +// newMetricsAdapter creates a metricsAdapter wrapping inner, registers all +// metrics on the given registry, and wires GaugeFunc callbacks via vol. +func newMetricsAdapter(inner iscsi.BlockDevice, vol *blockvol.BlockVol, reg prometheus.Registerer) *metricsAdapter { + writeOps := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "write_ops_total", Help: "Total write operations", + }) + readOps := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "read_ops_total", Help: "Total read operations", + }) + trimOps := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "trim_ops_total", Help: "Total trim operations", + }) + syncOps := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "sync_ops_total", Help: "Total sync operations", + }) + writeBytes := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "write_bytes_total", Help: "Total bytes written", + }) + readBytes := prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "read_bytes_total", Help: "Total bytes read", + }) + + latencyBuckets := []float64{0.00001, 0.00005, 0.0001, 0.0005, 0.001, 0.005, 0.01, 0.05, 0.1} + writeLatency := prometheus.NewHistogram(prometheus.HistogramOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "write_latency_seconds", Help: "Write latency distribution", + Buckets: latencyBuckets, + }) + readLatency := prometheus.NewHistogram(prometheus.HistogramOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "read_latency_seconds", Help: "Read latency distribution", + Buckets: latencyBuckets, + }) + syncLatency := prometheus.NewHistogram(prometheus.HistogramOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "sync_latency_seconds", Help: "Sync latency distribution", + Buckets: latencyBuckets, + }) + + reg.MustRegister(writeOps, readOps, trimOps, syncOps, writeBytes, readBytes, + writeLatency, readLatency, syncLatency) + + // Gauge callbacks from the engine. + gs := &gaugeSource{vol: vol} + reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "wal_used_fraction", Help: "WAL space usage (0.0 - 1.0)", + }, gs.walUsedFraction)) + reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "dirty_map_entries", Help: "Number of dirty map entries", + }, gs.dirtyMapEntries)) + reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "epoch", Help: "Current fencing epoch", + }, gs.epoch)) + reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "role", Help: "Current role (0=None, 1=Primary, 2=Replica, ...)", + }, gs.role)) + reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ + Namespace: "seaweedfs", Subsystem: "blockvol", + Name: "snapshot_count", Help: "Number of active snapshots", + }, gs.snapshotCount)) + + return &metricsAdapter{ + inner: inner, + writeOps: writeOps, + readOps: readOps, + trimOps: trimOps, + syncOps: syncOps, + writeBytes: writeBytes, + readBytes: readBytes, + writeLatency: writeLatency, + readLatency: readLatency, + syncLatency: syncLatency, + } +} + +func (m *metricsAdapter) ReadAt(lba uint64, length uint32) ([]byte, error) { + start := time.Now() + data, err := m.inner.ReadAt(lba, length) + m.readLatency.Observe(time.Since(start).Seconds()) + m.readOps.Inc() + m.readBytes.Add(float64(length)) + return data, err +} + +func (m *metricsAdapter) WriteAt(lba uint64, data []byte) error { + start := time.Now() + err := m.inner.WriteAt(lba, data) + m.writeLatency.Observe(time.Since(start).Seconds()) + m.writeOps.Inc() + m.writeBytes.Add(float64(len(data))) + return err +} + +func (m *metricsAdapter) Trim(lba uint64, length uint32) error { + err := m.inner.Trim(lba, length) + m.trimOps.Inc() + return err +} + +func (m *metricsAdapter) SyncCache() error { + start := time.Now() + err := m.inner.SyncCache() + m.syncLatency.Observe(time.Since(start).Seconds()) + m.syncOps.Inc() + return err +} + +func (m *metricsAdapter) BlockSize() uint32 { return m.inner.BlockSize() } +func (m *metricsAdapter) VolumeSize() uint64 { return m.inner.VolumeSize() } +func (m *metricsAdapter) IsHealthy() bool { return m.inner.IsHealthy() } + +// ALUAProvider proxy: delegate to inner device if it implements ALUAProvider. +func (m *metricsAdapter) ALUAState() uint8 { + if p, ok := m.inner.(iscsi.ALUAProvider); ok { + return p.ALUAState() + } + return iscsi.ALUAStandby +} +func (m *metricsAdapter) TPGroupID() uint16 { + if p, ok := m.inner.(iscsi.ALUAProvider); ok { + return p.TPGroupID() + } + return 1 +} +func (m *metricsAdapter) DeviceNAA() [8]byte { + if p, ok := m.inner.(iscsi.ALUAProvider); ok { + return p.DeviceNAA() + } + return [8]byte{} +} + +// --- gaugeSource callbacks --- + +func (gs *gaugeSource) walUsedFraction() float64 { + return gs.vol.WALUsedFraction() +} + +func (gs *gaugeSource) dirtyMapEntries() float64 { + return float64(gs.vol.DirtyMapLen()) +} + +func (gs *gaugeSource) epoch() float64 { + return float64(gs.vol.Status().Epoch) +} + +func (gs *gaugeSource) role() float64 { + return float64(gs.vol.Role()) +} + +func (gs *gaugeSource) snapshotCount() float64 { + return float64(len(gs.vol.ListSnapshots())) +} diff --git a/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics_test.go b/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics_test.go new file mode 100644 index 000000000..4e0f624d5 --- /dev/null +++ b/weed/storage/blockvol/iscsi/cmd/iscsi-target/metrics_test.go @@ -0,0 +1,110 @@ +package main + +import ( + "io" + "log" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/prometheus/client_golang/prometheus" + "github.com/seaweedfs/seaweedfs/weed/storage/blockvol" +) + +// TestMetrics_WriteIncrementsCounter writes 10 blocks and verifies +// write_ops_total >= 10 in the Prometheus output. +func TestMetrics_WriteIncrementsCounter(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "metrics.blk") + + vol, err := blockvol.CreateBlockVol(path, blockvol.CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + reg := prometheus.NewRegistry() + inner := &blockVolAdapter{vol: vol, tpgID: 1} + m := newMetricsAdapter(inner, vol, reg) + + // Write 10 blocks. + data := make([]byte, 4096) + for i := 0; i < 10; i++ { + data[0] = byte(i) + if err := m.WriteAt(uint64(i), data); err != nil { + t.Fatalf("write %d: %v", i, err) + } + } + + // Gather metrics. + families, err := reg.Gather() + if err != nil { + t.Fatalf("gather: %v", err) + } + found := false + for _, f := range families { + if f.GetName() == "seaweedfs_blockvol_write_ops_total" { + val := f.GetMetric()[0].GetCounter().GetValue() + if val < 10 { + t.Fatalf("write_ops_total: expected >= 10, got %v", val) + } + found = true + break + } + } + if !found { + t.Fatal("seaweedfs_blockvol_write_ops_total not found in metrics") + } +} + +// TestMetrics_EndpointServes starts the admin server and verifies that +// GET /metrics returns 200 with prometheus text format. +func TestMetrics_EndpointServes(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "endpoint.blk") + + vol, err := blockvol.CreateBlockVol(path, blockvol.CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + reg := prometheus.NewRegistry() + inner := &blockVolAdapter{vol: vol, tpgID: 1} + _ = newMetricsAdapter(inner, vol, reg) + + adm := newAdminServer(vol, "", log.New(os.Stderr, "[test] ", 0)) + adm.metricsRegistry = reg + ln, err := startAdminServer("127.0.0.1:0", adm) + if err != nil { + t.Fatalf("start admin: %v", err) + } + defer ln.Close() + + resp, err := http.Get("http://" + ln.Addr().String() + "/metrics") + if err != nil { + t.Fatalf("GET /metrics: %v", err) + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + t.Fatalf("expected 200, got %d", resp.StatusCode) + } + body, _ := io.ReadAll(resp.Body) + if !strings.Contains(string(body), "seaweedfs_blockvol_write_ops_total") { + t.Fatal("expected write_ops_total in metrics output") + } + if !strings.Contains(string(body), "seaweedfs_blockvol_wal_used_fraction") { + t.Fatal("expected wal_used_fraction in metrics output") + } +} diff --git a/weed/storage/blockvol/iscsi/login.go b/weed/storage/blockvol/iscsi/login.go index 19062915d..f2cd7b405 100644 --- a/weed/storage/blockvol/iscsi/login.go +++ b/weed/storage/blockvol/iscsi/login.go @@ -3,6 +3,7 @@ package iscsi import ( "errors" "strconv" + "strings" "time" ) @@ -72,6 +73,7 @@ type TargetConfig struct { ImmediateData bool ErrorRecoveryLevel int DataOutTimeout time.Duration // read deadline for Data-Out collection (default 30s) + CHAPConfig CHAPConfig // CHAP authentication settings } // DefaultTargetConfig returns sensible defaults for a target. @@ -101,6 +103,9 @@ type LoginNegotiator struct { tsih uint16 targetOK bool // target name validated + // CHAP authentication (nil when disabled) + chapAuth *CHAPAuthenticator + // Negotiated values (updated during negotiation) NegMaxRecvDataSegLen int NegMaxBurstLength int @@ -116,7 +121,7 @@ type LoginNegotiator struct { // NewLoginNegotiator creates a negotiator for a new login sequence. func NewLoginNegotiator(config TargetConfig) *LoginNegotiator { - return &LoginNegotiator{ + ln := &LoginNegotiator{ config: config, phase: LoginPhaseStart, NegMaxRecvDataSegLen: config.MaxRecvDataSegmentLength, @@ -125,6 +130,10 @@ func NewLoginNegotiator(config TargetConfig) *LoginNegotiator { NegInitialR2T: config.InitialR2T, NegImmediateData: config.ImmediateData, } + if config.CHAPConfig.Enabled { + ln.chapAuth = NewCHAPAuthenticator(config.CHAPConfig) + } + return ln } // HandleLoginPDU processes one login request PDU and returns the response PDU. @@ -198,8 +207,65 @@ func (ln *LoginNegotiator) HandleLoginPDU(req *PDU, resolver TargetResolver) *PD // ISID ln.isid = req.ISID() - // We don't implement CHAP -- declare AuthMethod=None - respParams.Set("AuthMethod", "None") + // CHAP authentication flow + if ln.chapAuth != nil && ln.chapAuth.IsEnabled() { + authMethod, _ := params.Get("AuthMethod") + chapN, hasChapN := params.Get("CHAP_N") + chapR, hasChapR := params.Get("CHAP_R") + + switch ln.chapAuth.state { + case chapIdle: + // First security PDU: initiator offers AuthMethod. + // Check if initiator supports CHAP. + if !chapMethodOffered(authMethod) { + // CHAP required but initiator only offers None. + setLoginReject(resp, LoginStatusInitiatorErr, LoginDetailAuthFailure) + return resp + } + respParams.Set("AuthMethod", "CHAP") + challenge, err := ln.chapAuth.GenerateChallenge() + if err != nil { + setLoginReject(resp, LoginStatusTargetErr, LoginDetailTargetError) + return resp + } + for k, v := range challenge { + respParams.Set(k, v) + } + // Do NOT transit yet -- more security PDUs needed. + resp.SetLoginStages(csg, nsg) + resp.SetLoginTransit(false) + resp.SetLoginStatus(LoginStatusSuccess, LoginDetailSuccess) + if ln.tsih == 0 { + ln.tsih = 1 + } + resp.SetTSIH(ln.tsih) + tpgt := ln.config.TargetPortalGroupTag + if tpgt <= 0 { + tpgt = 1 + } + respParams.Set("TargetPortalGroupTag", strconv.Itoa(tpgt)) + if respParams.Len() > 0 { + resp.DataSegment = respParams.Encode() + } + return resp + + case chapChallengeSent: + // Second security PDU: initiator sends CHAP_N + CHAP_R. + if !hasChapN || !hasChapR { + setLoginReject(resp, LoginStatusInitiatorErr, LoginDetailAuthFailure) + return resp + } + if !ln.chapAuth.Verify(chapN, chapR) { + setLoginReject(resp, LoginStatusInitiatorErr, LoginDetailAuthFailure) + return resp + } + // CHAP verified -- echo AuthMethod, allow transit below. + respParams.Set("AuthMethod", "CHAP") + } + } else { + // No CHAP -- declare AuthMethod=None. + respParams.Set("AuthMethod", "None") + } if transit { if nsg == StageLoginOp { @@ -376,6 +442,17 @@ func setLoginReject(resp *PDU, class, detail uint8) { resp.SetLoginTransit(false) } +// chapMethodOffered checks whether "CHAP" appears in a comma-separated +// AuthMethod value list (e.g. "CHAP,None" or "CHAP"). +func chapMethodOffered(authMethod string) bool { + for _, m := range strings.Split(authMethod, ",") { + if strings.TrimSpace(m) == "CHAP" { + return true + } + } + return false +} + // LoginResult contains the outcome of a completed login negotiation. type LoginResult struct { InitiatorName string diff --git a/weed/storage/blockvol/resize_test.go b/weed/storage/blockvol/resize_test.go new file mode 100644 index 000000000..ac3750ff7 --- /dev/null +++ b/weed/storage/blockvol/resize_test.go @@ -0,0 +1,144 @@ +package blockvol + +import ( + "bytes" + "os" + "path/filepath" + "testing" +) + +// TestResize_ExpandWorks grows a 1MB volume to 2MB, writes to the new region, +// and reads back to verify. +func TestResize_ExpandWorks(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "expand.blk") + + vol, err := CreateBlockVol(path, CreateOptions{ + VolumeSize: 1024 * 1024, // 1MB + BlockSize: 4096, + WALSize: 64 * 1024, // 64KB (small for test) + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + // Write to last block of original size. + lastLBA := uint64((1024*1024)/4096 - 1) + data := make([]byte, 4096) + for i := range data { + data[i] = 0xAA + } + if err := vol.WriteLBA(lastLBA, data); err != nil { + t.Fatalf("write last block: %v", err) + } + + // Expand to 2MB. + newSize := uint64(2 * 1024 * 1024) + if err := vol.Expand(newSize); err != nil { + t.Fatalf("expand: %v", err) + } + + // Verify volume size updated. + if vol.Info().VolumeSize != newSize { + t.Fatalf("expected VolumeSize=%d, got %d", newSize, vol.Info().VolumeSize) + } + + // Write to a block in the new region. + newLBA := uint64((1024 * 1024) / 4096) // first block in expanded region + data2 := make([]byte, 4096) + for i := range data2 { + data2[i] = 0xBB + } + if err := vol.WriteLBA(newLBA, data2); err != nil { + t.Fatalf("write new region: %v", err) + } + + // Read back old data. + got, err := vol.ReadLBA(lastLBA, 4096) + if err != nil { + t.Fatalf("read old block: %v", err) + } + if !bytes.Equal(got, data) { + t.Fatal("old block data mismatch after expand") + } + + // Read back new data. + got2, err := vol.ReadLBA(newLBA, 4096) + if err != nil { + t.Fatalf("read new block: %v", err) + } + if !bytes.Equal(got2, data2) { + t.Fatal("new block data mismatch after expand") + } + + // Verify file size on disk. + fi, _ := os.Stat(path) + extentStart := vol.super.WALOffset + vol.super.WALSize + expected := int64(extentStart + newSize) + if fi.Size() != expected { + t.Fatalf("file size: expected %d, got %d", expected, fi.Size()) + } +} + +// TestResize_ShrinkRejected verifies that attempting to shrink returns an error. +func TestResize_ShrinkRejected(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "shrink.blk") + + vol, err := CreateBlockVol(path, CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + err = vol.Expand(512 * 1024) // shrink + if err != ErrShrinkNotSupported { + t.Fatalf("expected ErrShrinkNotSupported, got %v", err) + } + + // Same size = no-op. + if err := vol.Expand(1024 * 1024); err != nil { + t.Fatalf("same-size expand should be no-op: %v", err) + } +} + +// TestResize_WithSnapshotsRejected verifies that resize is blocked when +// snapshots are active. +func TestResize_WithSnapshotsRejected(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "snap-resize.blk") + + vol, err := CreateBlockVol(path, CreateOptions{ + VolumeSize: 1024 * 1024, + BlockSize: 4096, + WALSize: 64 * 1024, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + defer vol.Close() + + // Create a snapshot. + if err := vol.CreateSnapshot(1); err != nil { + t.Fatalf("create snapshot: %v", err) + } + + // Try to expand -- should fail. + err = vol.Expand(2 * 1024 * 1024) + if err != ErrSnapshotsPreventResize { + t.Fatalf("expected ErrSnapshotsPreventResize, got %v", err) + } + + // Delete snapshot, then expand should work. + if err := vol.DeleteSnapshot(1); err != nil { + t.Fatalf("delete snapshot: %v", err) + } + if err := vol.Expand(2 * 1024 * 1024); err != nil { + t.Fatalf("expand after snapshot delete: %v", err) + } +}