(ctx ModuleContext, sessionID string, taskID uint32, shellTool string, sess *sessions.Session, chunks []sessions.UploadChunk)
| 88 | } |
| 89 | |
| 90 | func (m *UploadModule) executeChunks(ctx ModuleContext, sessionID string, taskID uint32, shellTool string, sess *sessions.Session, chunks []sessions.UploadChunk) { |
| 91 | ch := ctx.Tasks.AwaitResult(sessionID, taskID) |
| 92 | if ch == nil { |
| 93 | ctx.Tasks.Fail(sessionID, taskID, "await failed") |
| 94 | return |
| 95 | } |
| 96 | |
| 97 | for i, chunk := range chunks { |
| 98 | args := sessions.BuildCommandArguments(sess, shellTool, chunk.Command) |
| 99 | if _, ok := enqueueToolAction(ctx, sessionID, taskID, shellTool, args); !ok { |
| 100 | sendUploadACK(ctx, sessionID, taskID, false) |
| 101 | return |
| 102 | } |
| 103 | log.Infof("[bridge] enqueued upload chunk %d/%d for session %s", i+1, len(chunks), sessionID) |
| 104 | |
| 105 | // Wait for this chunk's result before enqueuing the next. |
| 106 | if _, ok := awaitTaskResult(ch, taskID); !ok { |
| 107 | log.Errorf("[bridge] upload chunk %d/%d failed for session %s", i+1, len(chunks), sessionID) |
| 108 | sendUploadACK(ctx, sessionID, taskID, false) |
| 109 | ctx.Tasks.Fail(sessionID, taskID, "chunk failed") |
| 110 | return |
| 111 | } |
| 112 | } |
| 113 | |
| 114 | log.Infof("[bridge] all %d upload chunks completed for session %s", len(chunks), sessionID) |
| 115 | sendUploadACK(ctx, sessionID, taskID, true) |
| 116 | ctx.Tasks.Complete(sessionID, taskID) |
| 117 | } |
| 118 | |
| 119 | // sendUploadACK sends an upload acknowledgment via ModuleContext. |
| 120 | func sendUploadACK(ctx ModuleContext, sessionID string, taskID uint32, success bool) { |
no test coverage detected