mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Feat/deployment hostpath (#24985)
* feat(llm): hostpath support in deployment * fix(llm): status error in deployment & sku deleted by deployment accidentally
This commit is contained in:
@@ -113,6 +113,8 @@ type LLMDeploymentCreateInput struct {
|
||||
AutoStart bool `json:"auto_start"`
|
||||
// Prefer specific host
|
||||
PreferHost string `json:"prefer_host"`
|
||||
// Host path mounts for instances.
|
||||
HostPaths *HostPaths `json:"host_paths,omitempty"`
|
||||
|
||||
// Expected number of replicas (SLLM instances)
|
||||
Replicas int `json:"replicas"`
|
||||
|
||||
+14
-9
@@ -488,11 +488,9 @@ func (llm *SLLM) CustomizeDelete(ctx context.Context, userCred mcclient.TokenCre
|
||||
return llm.StartDeleteTask(ctx, userCred, "")
|
||||
}
|
||||
|
||||
// SetStatus overrides the base implementation to push ReadyReplicas updates to
|
||||
// the parent SLLMDeployment whenever the instance crosses the running boundary.
|
||||
// This is the only deterministic hook for "instance came up" / "instance went
|
||||
// down" — every transition (create-complete, start, stop, fail, restart) flows
|
||||
// through here.
|
||||
// SetStatus overrides the base implementation to push replica health updates
|
||||
// to the parent SLLMDeployment whenever the instance crosses the running
|
||||
// boundary or enters/leaves a deployment-relevant failure status.
|
||||
func (llm *SLLM) SetStatus(ctx context.Context, userCred mcclient.TokenCredential, status string, reason string) error {
|
||||
oldStatus := llm.Status
|
||||
if err := llm.SLLMBase.SetStatus(ctx, userCred, status, reason); err != nil {
|
||||
@@ -501,10 +499,7 @@ func (llm *SLLM) SetStatus(ctx context.Context, userCred mcclient.TokenCredentia
|
||||
if llm.LLMDeploymentId == "" {
|
||||
return nil
|
||||
}
|
||||
// Only sync when running-membership changed.
|
||||
wasRunning := oldStatus == api.LLM_STATUS_RUNNING
|
||||
isRunning := status == api.LLM_STATUS_RUNNING
|
||||
if wasRunning == isRunning {
|
||||
if !shouldSyncDeploymentOnLLMStatusChange(oldStatus, status) {
|
||||
return nil
|
||||
}
|
||||
depObj, err := GetLLMDeploymentManager().FetchById(llm.LLMDeploymentId)
|
||||
@@ -518,6 +513,16 @@ func (llm *SLLM) SetStatus(ctx context.Context, userCred mcclient.TokenCredentia
|
||||
return nil
|
||||
}
|
||||
|
||||
func shouldSyncDeploymentOnLLMStatusChange(oldStatus string, newStatus string) bool {
|
||||
if oldStatus == newStatus {
|
||||
return false
|
||||
}
|
||||
if oldStatus == api.LLM_STATUS_RUNNING || newStatus == api.LLM_STATUS_RUNNING {
|
||||
return true
|
||||
}
|
||||
return isDeploymentReplicaFailureStatus(oldStatus) != isDeploymentReplicaFailureStatus(newStatus)
|
||||
}
|
||||
|
||||
func (llm *SLLM) GetLLMSku(skuId string) (*SLLMSku, error) {
|
||||
if len(skuId) == 0 {
|
||||
skuId = llm.LLMSkuId
|
||||
|
||||
@@ -53,6 +53,10 @@ type SLLMDeployment struct {
|
||||
// Nullable: may be empty during the brief window when a deployment auto-creates a SKU via SkuSpec.
|
||||
LLMSkuId string `width:"128" charset:"ascii" nullable:"true" list:"user" create:"optional"`
|
||||
|
||||
// SKU automatically created for this deployment and owned by its lifecycle.
|
||||
// Empty when the deployment reuses an existing SKU.
|
||||
ManagedLLMSkuId string `width:"128" charset:"ascii" nullable:"true" list:"user"`
|
||||
|
||||
// Expected replica count
|
||||
Replicas int `nullable:"false" default:"1" list:"user" create:"optional" update:"user"`
|
||||
|
||||
@@ -93,10 +97,11 @@ type SLLMDeployment struct {
|
||||
AccessPolicy string `width:"32" charset:"ascii" nullable:"true" default:"authed" list:"user" create:"optional" update:"user"`
|
||||
|
||||
// Instance template — captured at create time so scale-up can re-create
|
||||
// new SLLM replicas with the same network / start settings as the originals.
|
||||
// new SLLM replicas with the same network / start / mount settings as the originals.
|
||||
// Scale request bodies don't carry these fields.
|
||||
Nets *api.LLMDeploymentNets `charset:"utf8" length:"medium" nullable:"true" list:"user"`
|
||||
AutoStart bool `nullable:"false" default:"false" list:"user"`
|
||||
HostPaths *api.HostPaths `charset:"utf8" length:"medium" nullable:"true" list:"user"`
|
||||
}
|
||||
|
||||
func (man *SLLMDeploymentManager) ValidateCreateData(
|
||||
@@ -373,71 +378,103 @@ func (model *SLLMDeployment) RealDelete(ctx context.Context, userCred mcclient.T
|
||||
return model.SVirtualResourceBase.Delete(ctx, userCred)
|
||||
}
|
||||
|
||||
// SyncReadyReplicas recomputes ReadyReplicas from running SLLM instances,
|
||||
// persists it to the deployment row, and transitions the deployment status
|
||||
// based on how many replicas are up:
|
||||
// SyncReadyReplicas recomputes ReadyReplicas from SLLM instances, persists it
|
||||
// to the deployment row, and transitions the deployment status based on replica
|
||||
// health:
|
||||
//
|
||||
// ready_replicas == 0 && desired > 0 → deploying
|
||||
// 0 < ready_replicas < desired → partial
|
||||
// ready_replicas == desired → ready
|
||||
// replica create/start failure → create_fail (unless some replicas run)
|
||||
//
|
||||
// Failure / lifecycle statuses (create_fail, deleting, importing_model, etc.)
|
||||
// are not overridden — the running-state status only takes effect once the
|
||||
// deployment has cleared the early lifecycle and has at least one instance.
|
||||
// are not overridden once set.
|
||||
//
|
||||
// Call after create/scale tasks finish and on every instance status change.
|
||||
func (model *SLLMDeployment) SyncReadyReplicas(ctx context.Context, userCred mcclient.TokenCredential) error {
|
||||
cnt, err := GetLLMManager().Query().
|
||||
var rows []deploymentReplicaStatusRow
|
||||
err := GetLLMManager().Query("status").
|
||||
Equals("llm_deployment_id", model.Id).
|
||||
Equals("status", api.LLM_STATUS_RUNNING).
|
||||
CountWithError()
|
||||
All(&rows)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "count running instances")
|
||||
return errors.Wrap(err, "fetch replica statuses")
|
||||
}
|
||||
if model.ReadyReplicas != cnt {
|
||||
summary := summarizeDeploymentReplicaStatuses(rows)
|
||||
if model.ReadyReplicas != summary.Running {
|
||||
if _, err = db.Update(model, func() error {
|
||||
model.ReadyReplicas = cnt
|
||||
model.ReadyReplicas = summary.Running
|
||||
return nil
|
||||
}); err != nil {
|
||||
return errors.Wrap(err, "update ready_replicas")
|
||||
}
|
||||
}
|
||||
|
||||
// Map (ready_replicas, replicas) → desired deployment status.
|
||||
desired := computeRunningStatus(cnt, model.Replicas)
|
||||
desired := computeDeploymentReplicaStatus(summary, model.Replicas)
|
||||
if desired == "" {
|
||||
return nil
|
||||
}
|
||||
if !canFlipToRunningStatus(model.Status) {
|
||||
if !canUpdateReplicaHealthStatus(model.Status) {
|
||||
return nil
|
||||
}
|
||||
if model.Status == desired {
|
||||
return nil
|
||||
}
|
||||
return model.SetStatus(ctx, userCred, desired, fmt.Sprintf("ready_replicas=%d/%d", cnt, model.Replicas))
|
||||
return model.SetStatus(ctx, userCred, desired, fmt.Sprintf("ready_replicas=%d/%d", summary.Running, model.Replicas))
|
||||
}
|
||||
|
||||
// computeRunningStatus returns the status that reflects the running replica
|
||||
// count. Returns "" when there's nothing to set (e.g., desired == 0).
|
||||
func computeRunningStatus(ready, desired int) string {
|
||||
type deploymentReplicaStatusRow struct {
|
||||
Status string `json:"status"`
|
||||
}
|
||||
|
||||
type deploymentReplicaStatusSummary struct {
|
||||
Running int
|
||||
HasFailure bool
|
||||
}
|
||||
|
||||
func summarizeDeploymentReplicaStatuses(rows []deploymentReplicaStatusRow) deploymentReplicaStatusSummary {
|
||||
summary := deploymentReplicaStatusSummary{}
|
||||
for i := range rows {
|
||||
switch rows[i].Status {
|
||||
case api.LLM_STATUS_RUNNING:
|
||||
summary.Running++
|
||||
default:
|
||||
if isDeploymentReplicaFailureStatus(rows[i].Status) {
|
||||
summary.HasFailure = true
|
||||
}
|
||||
}
|
||||
}
|
||||
return summary
|
||||
}
|
||||
|
||||
func computeDeploymentReplicaStatus(summary deploymentReplicaStatusSummary, desired int) string {
|
||||
if desired <= 0 {
|
||||
return ""
|
||||
}
|
||||
switch {
|
||||
case ready == 0:
|
||||
return api.LLM_DEPLOYMENT_STATUS_DEPLOYING
|
||||
case ready < desired:
|
||||
return api.LLM_DEPLOYMENT_STATUS_PARTIAL
|
||||
default:
|
||||
case summary.Running >= desired:
|
||||
return api.STATUS_READY
|
||||
case summary.Running > 0:
|
||||
return api.LLM_DEPLOYMENT_STATUS_PARTIAL
|
||||
case summary.HasFailure:
|
||||
return api.LLM_STATUS_CREATE_FAIL
|
||||
default:
|
||||
return api.LLM_DEPLOYMENT_STATUS_DEPLOYING
|
||||
}
|
||||
}
|
||||
|
||||
// canFlipToRunningStatus reports whether the deployment is in a phase where
|
||||
// the running-replica-driven status can override the current value. Early
|
||||
// lifecycle (importing_model, creating_sku) and terminal failure / delete
|
||||
// states must not be clobbered.
|
||||
func canFlipToRunningStatus(current string) bool {
|
||||
func isDeploymentReplicaFailureStatus(status string) bool {
|
||||
switch status {
|
||||
case api.LLM_STATUS_CREATE_FAIL,
|
||||
api.LLM_STATUS_START_FAIL:
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// canUpdateReplicaHealthStatus reports whether replica-health-driven status can
|
||||
// override the current deployment value. Early lifecycle and terminal failure /
|
||||
// delete states must not be clobbered.
|
||||
func canUpdateReplicaHealthStatus(current string) bool {
|
||||
switch current {
|
||||
case api.LLM_DEPLOYMENT_STATUS_IMPORTING_MODEL,
|
||||
api.LLM_DEPLOYMENT_STATUS_IMPORT_MODEL_FAILED,
|
||||
@@ -490,7 +527,7 @@ func TriggerLLMDeploymentReconcile(ctx context.Context, userCred mcclient.TokenC
|
||||
func (model *SLLMDeployment) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) {
|
||||
model.SVirtualResourceBase.PostCreate(ctx, userCred, ownerId, query, data)
|
||||
|
||||
// Persist instance template (nets / auto_start) so later scale operations
|
||||
// Persist instance template (nets / auto_start / host_paths) so later scale operations
|
||||
// can rebuild instances without these fields in the scale request body.
|
||||
input := api.LLMDeploymentCreateInput{}
|
||||
_ = data.Unmarshal(&input)
|
||||
@@ -500,6 +537,9 @@ func (model *SLLMDeployment) PostCreate(ctx context.Context, userCred mcclient.T
|
||||
model.Nets = &nets
|
||||
}
|
||||
model.AutoStart = input.AutoStart
|
||||
if input.HostPaths != nil && !input.HostPaths.IsZero() {
|
||||
model.HostPaths = input.HostPaths
|
||||
}
|
||||
return nil
|
||||
}); err != nil {
|
||||
log.Errorf("SLLMDeployment.PostCreate persist instance template: %s", err)
|
||||
|
||||
@@ -255,8 +255,7 @@ func (task *LLMDeploymentCreateTask) createSkuAndReconcile(ctx context.Context,
|
||||
}
|
||||
|
||||
if _, err := db.Update(model, func() error {
|
||||
model.LLMSkuId = skuId
|
||||
return nil
|
||||
return assignDeploymentAutoCreatedSku(model, skuId)
|
||||
}); err != nil {
|
||||
task.taskFailedCreatingSku(ctx, model, errors.Wrap(err, "write back llm_sku_id"))
|
||||
return
|
||||
@@ -266,6 +265,12 @@ func (task *LLMDeploymentCreateTask) createSkuAndReconcile(ctx context.Context,
|
||||
task.reconcileAndComplete(ctx, model, body)
|
||||
}
|
||||
|
||||
func assignDeploymentAutoCreatedSku(model *models.SLLMDeployment, skuId string) error {
|
||||
model.LLMSkuId = skuId
|
||||
model.ManagedLLMSkuId = skuId
|
||||
return nil
|
||||
}
|
||||
|
||||
// reconcileReplicas is the core logic that ensures actual SLLM count matches desired replicas.
|
||||
// - If actual < desired: create new instances (scale up)
|
||||
// - If actual > desired: delete excess instances (scale down, prefer unhealthy)
|
||||
@@ -333,7 +338,7 @@ func scaleUp(ctx context.Context, userCred mcclient.TokenCredential, deployment
|
||||
}
|
||||
}
|
||||
|
||||
// Instance template (nets / auto_start) lives on the deployment row,
|
||||
// Instance template (nets / auto_start / host_paths) lives on the deployment row,
|
||||
// captured at create time by SLLMDeployment.PostCreate. Scale request
|
||||
// bodies don't carry these fields, so the row is the source of truth.
|
||||
if deployment.Nets == nil || len(*deployment.Nets) == 0 {
|
||||
@@ -356,16 +361,7 @@ func scaleUp(ctx context.Context, userCred mcclient.TokenCredential, deployment
|
||||
instanceName := fmt.Sprintf("%s-%d", deployment.Name, nextIndex+i)
|
||||
|
||||
// Build typed LLMCreateInput, then marshal to JSON for handler.Create
|
||||
llmInput := api.LLMCreateInput{
|
||||
LLMBaseCreateInput: api.LLMBaseCreateInput{
|
||||
AutoStart: deployment.AutoStart,
|
||||
Nets: nets,
|
||||
},
|
||||
LLMSkuId: deployment.LLMSkuId,
|
||||
LLMImageId: imageId,
|
||||
LLMDeploymentId: deployment.Id,
|
||||
LLMSpec: llmSpec,
|
||||
}
|
||||
llmInput := buildDeploymentLLMCreateInput(deployment, nets, imageId, llmSpec)
|
||||
|
||||
llmParams := jsonutils.Marshal(llmInput).(*jsonutils.JSONDict)
|
||||
// Name lives on the embedded VirtualResourceCreateInput → set explicitly
|
||||
@@ -390,6 +386,20 @@ func scaleUp(ctx context.Context, userCred mcclient.TokenCredential, deployment
|
||||
return nil
|
||||
}
|
||||
|
||||
func buildDeploymentLLMCreateInput(deployment *models.SLLMDeployment, nets []*computeapi.NetworkConfig, imageId string, llmSpec *api.LLMSpec) api.LLMCreateInput {
|
||||
return api.LLMCreateInput{
|
||||
LLMBaseCreateInput: api.LLMBaseCreateInput{
|
||||
AutoStart: deployment.AutoStart,
|
||||
Nets: nets,
|
||||
},
|
||||
LLMSkuId: deployment.LLMSkuId,
|
||||
LLMImageId: imageId,
|
||||
LLMDeploymentId: deployment.Id,
|
||||
LLMSpec: llmSpec,
|
||||
HostPaths: deployment.HostPaths,
|
||||
}
|
||||
}
|
||||
|
||||
// scaleDown deletes `count` SLLM instances, preferring unhealthy ones.
|
||||
func scaleDown(ctx context.Context, userCred mcclient.TokenCredential, model *models.SLLMDeployment, instances []instanceInfo, count int) error {
|
||||
// Sort: unhealthy/error/stopped first, then by name descending (newest first)
|
||||
|
||||
@@ -96,8 +96,9 @@ func (task *LLMDeploymentDeleteTask) OnInstancesDeletedFailed(ctx context.Contex
|
||||
}
|
||||
|
||||
func (task *LLMDeploymentDeleteTask) deleteDeployment(ctx context.Context, model *models.SLLMDeployment) {
|
||||
// Capture associated SKU id BEFORE deleting the deployment so we can cascade.
|
||||
skuId := model.LLMSkuId
|
||||
// Capture managed SKU id BEFORE deleting the deployment so we can cascade
|
||||
// only SKUs automatically created for this deployment.
|
||||
skuId := deploymentManagedSkuIdForCascade(model)
|
||||
|
||||
if err := model.RealDelete(ctx, task.UserCred); err != nil {
|
||||
log.Errorf("LLMDeploymentDeleteTask: RealDelete deployment %s: %s", model.Id, err)
|
||||
@@ -113,6 +114,13 @@ func (task *LLMDeploymentDeleteTask) deleteDeployment(ctx context.Context, model
|
||||
task.SetStageComplete(ctx, nil)
|
||||
}
|
||||
|
||||
func deploymentManagedSkuIdForCascade(model *models.SLLMDeployment) string {
|
||||
if model == nil {
|
||||
return ""
|
||||
}
|
||||
return model.ManagedLLMSkuId
|
||||
}
|
||||
|
||||
// cascadeDeleteSku tries to delete the SKU created for this deployment. If the
|
||||
// SKU is still referenced (e.g., by another deployment's instances), the
|
||||
// framework's ValidateDeleteCondition rejects it and we log + continue.
|
||||
|
||||
@@ -18,7 +18,7 @@ import (
|
||||
// desired value.
|
||||
//
|
||||
// Unlike LLMDeploymentCreateTask, this task's body does NOT carry the original
|
||||
// create-time payload (nets / auto_start / prefer_host). Those fields are read
|
||||
// create-time payload (nets / auto_start / host_paths). Those fields are read
|
||||
// from the deployment row, where PostCreate persisted them.
|
||||
type LLMDeploymentSyncReplicasTask struct {
|
||||
taskman.STask
|
||||
|
||||
@@ -74,6 +74,7 @@ type LLMDeploymentCreateOptions struct {
|
||||
Net []string `help:"Network descriptions; repeatable" json:"-"`
|
||||
AutoStart bool `help:"auto start instances after creation" json:"auto_start"`
|
||||
PreferHost string `help:"prefer specific host" json:"prefer_host"`
|
||||
HostPaths []string `json:"-" help:"host path mount in format path=<host_path>,type=<directory|file>,container_index=<index>,mount_path=<container_path>[,auto_create=<bool>][,read_only=<bool>][,propagation=<private|rslave|rshared>][,fs_user=<uid>][,fs_group=<gid>][,uid=<uid>][,gid=<gid>][,permissions=<mode>]; repeatable"`
|
||||
|
||||
// Deployment config
|
||||
Replicas int `help:"number of replicas" default:"1" json:"replicas"`
|
||||
@@ -125,6 +126,9 @@ func (o *LLMDeploymentCreateOptions) Params() (jsonutils.JSONObject, error) {
|
||||
}
|
||||
params.Set("nets", jsonutils.Marshal(nets))
|
||||
}
|
||||
if err := fetchHostPaths(o.HostPaths, params); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return params, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user