You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

438 lines
14 KiB

7 years ago
7 years ago
7 years ago
7 years ago
4 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
4 years ago
4 years ago
4 years ago
  1. package s3api
  2. import (
  3. "context"
  4. "encoding/xml"
  5. "fmt"
  6. "github.com/seaweedfs/seaweedfs/weed/glog"
  7. "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
  8. "io"
  9. "net/http"
  10. "net/url"
  11. "path/filepath"
  12. "strconv"
  13. "strings"
  14. "time"
  15. "github.com/seaweedfs/seaweedfs/weed/filer"
  16. "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
  17. "github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
  18. )
  19. type ListBucketResultV2 struct {
  20. XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ ListBucketResult"`
  21. Name string `xml:"Name"`
  22. Prefix string `xml:"Prefix"`
  23. MaxKeys int `xml:"MaxKeys"`
  24. Delimiter string `xml:"Delimiter,omitempty"`
  25. IsTruncated bool `xml:"IsTruncated"`
  26. Contents []ListEntry `xml:"Contents,omitempty"`
  27. CommonPrefixes []PrefixEntry `xml:"CommonPrefixes,omitempty"`
  28. ContinuationToken string `xml:"ContinuationToken,omitempty"`
  29. NextContinuationToken string `xml:"NextContinuationToken,omitempty"`
  30. KeyCount int `xml:"KeyCount"`
  31. StartAfter string `xml:"StartAfter,omitempty"`
  32. }
  33. func (s3a *S3ApiServer) ListObjectsV2Handler(w http.ResponseWriter, r *http.Request) {
  34. // https://docs.aws.amazon.com/AmazonS3/latest/API/v2-RESTBucketGET.html
  35. // collect parameters
  36. bucket, _ := s3_constants.GetBucketAndObject(r)
  37. glog.V(3).Infof("ListObjectsV2Handler %s", bucket)
  38. originalPrefix, continuationToken, startAfter, delimiter, _, maxKeys := getListObjectsV2Args(r.URL.Query())
  39. if maxKeys < 0 {
  40. s3err.WriteErrorResponse(w, r, s3err.ErrInvalidMaxKeys)
  41. return
  42. }
  43. if delimiter != "" && delimiter != "/" {
  44. s3err.WriteErrorResponse(w, r, s3err.ErrNotImplemented)
  45. return
  46. }
  47. marker := continuationToken
  48. if continuationToken == "" {
  49. marker = startAfter
  50. }
  51. response, err := s3a.listFilerEntries(bucket, originalPrefix, maxKeys, marker, delimiter)
  52. if err != nil {
  53. s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
  54. return
  55. }
  56. if len(response.Contents) == 0 {
  57. if exists, existErr := s3a.exists(s3a.option.BucketsPath, bucket, true); existErr == nil && !exists {
  58. s3err.WriteErrorResponse(w, r, s3err.ErrNoSuchBucket)
  59. return
  60. }
  61. }
  62. responseV2 := &ListBucketResultV2{
  63. XMLName: response.XMLName,
  64. Name: response.Name,
  65. CommonPrefixes: response.CommonPrefixes,
  66. Contents: response.Contents,
  67. ContinuationToken: continuationToken,
  68. Delimiter: response.Delimiter,
  69. IsTruncated: response.IsTruncated,
  70. KeyCount: len(response.Contents) + len(response.CommonPrefixes),
  71. MaxKeys: response.MaxKeys,
  72. NextContinuationToken: response.NextMarker,
  73. Prefix: response.Prefix,
  74. StartAfter: startAfter,
  75. }
  76. writeSuccessResponseXML(w, r, responseV2)
  77. }
  78. func (s3a *S3ApiServer) ListObjectsV1Handler(w http.ResponseWriter, r *http.Request) {
  79. // https://docs.aws.amazon.com/AmazonS3/latest/API/RESTBucketGET.html
  80. // collect parameters
  81. bucket, _ := s3_constants.GetBucketAndObject(r)
  82. glog.V(3).Infof("ListObjectsV1Handler %s", bucket)
  83. originalPrefix, marker, delimiter, maxKeys := getListObjectsV1Args(r.URL.Query())
  84. if maxKeys < 0 {
  85. s3err.WriteErrorResponse(w, r, s3err.ErrInvalidMaxKeys)
  86. return
  87. }
  88. if delimiter != "" && delimiter != "/" {
  89. s3err.WriteErrorResponse(w, r, s3err.ErrNotImplemented)
  90. return
  91. }
  92. response, err := s3a.listFilerEntries(bucket, originalPrefix, maxKeys, marker, delimiter)
  93. if err != nil {
  94. s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
  95. return
  96. }
  97. if len(response.Contents) == 0 {
  98. if exists, existErr := s3a.exists(s3a.option.BucketsPath, bucket, true); existErr == nil && !exists {
  99. s3err.WriteErrorResponse(w, r, s3err.ErrNoSuchBucket)
  100. return
  101. }
  102. }
  103. writeSuccessResponseXML(w, r, response)
  104. }
  105. func (s3a *S3ApiServer) listFilerEntries(bucket string, originalPrefix string, maxKeys int, marker string, delimiter string) (response ListBucketResult, err error) {
  106. // convert full path prefix into directory name and prefix for entry name
  107. reqDir, prefix := filepath.Split(originalPrefix)
  108. if strings.HasPrefix(reqDir, "/") {
  109. reqDir = reqDir[1:]
  110. }
  111. bucketPrefix := fmt.Sprintf("%s/%s/", s3a.option.BucketsPath, bucket)
  112. bucketPrefixLen := len(bucketPrefix)
  113. reqDir = fmt.Sprintf("%s%s", bucketPrefix, reqDir)
  114. if strings.HasSuffix(reqDir, "/") {
  115. reqDir = strings.TrimSuffix(reqDir, "/")
  116. }
  117. var contents []ListEntry
  118. var commonPrefixes []PrefixEntry
  119. var isTruncated bool
  120. var doErr error
  121. var nextMarker string
  122. // check filer
  123. err = s3a.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
  124. _, isTruncated, nextMarker, doErr = s3a.doListFilerEntries(client, reqDir, prefix, maxKeys, marker, delimiter, false, false, bucketPrefixLen, func(dir string, entry *filer_pb.Entry) {
  125. if entry.IsDirectory {
  126. if delimiter == "/" {
  127. commonPrefixes = append(commonPrefixes, PrefixEntry{
  128. Prefix: fmt.Sprintf("%s/%s/", dir, entry.Name)[bucketPrefixLen:],
  129. })
  130. }
  131. if !(entry.IsDirectoryKeyObject() && strings.HasSuffix(entry.Name, "/")) {
  132. return
  133. }
  134. }
  135. storageClass := "STANDARD"
  136. if v, ok := entry.Extended[s3_constants.AmzStorageClass]; ok {
  137. storageClass = string(v)
  138. }
  139. contents = append(contents, ListEntry{
  140. Key: fmt.Sprintf("%s/%s", dir, entry.Name)[bucketPrefixLen:],
  141. LastModified: time.Unix(entry.Attributes.Mtime, 0).UTC(),
  142. ETag: "\"" + filer.ETag(entry) + "\"",
  143. Size: int64(filer.FileSize(entry)),
  144. Owner: CanonicalUser{
  145. ID: fmt.Sprintf("%x", entry.Attributes.Uid),
  146. DisplayName: entry.Attributes.UserName,
  147. },
  148. StorageClass: StorageClass(storageClass),
  149. })
  150. })
  151. glog.V(4).Infof("end doListFilerEntries isTruncated:%v nextMarker:%v reqDir: %v prefix: %v", isTruncated, nextMarker, reqDir, prefix)
  152. if doErr != nil {
  153. return doErr
  154. }
  155. if !isTruncated {
  156. nextMarker = ""
  157. }
  158. if len(contents) == 0 && len(commonPrefixes) == 0 && maxKeys > 0 {
  159. if strings.HasSuffix(originalPrefix, "/") && prefix == "" {
  160. reqDir, prefix = filepath.Split(strings.TrimSuffix(reqDir, "/"))
  161. reqDir = strings.TrimSuffix(reqDir, "/")
  162. }
  163. _, _, _, doErr = s3a.doListFilerEntries(client, reqDir, prefix, 1, prefix, delimiter, true, false, bucketPrefixLen, func(dir string, entry *filer_pb.Entry) {
  164. if entry.IsDirectoryKeyObject() && entry.Name == prefix {
  165. storageClass := "STANDARD"
  166. if v, ok := entry.Extended[s3_constants.AmzStorageClass]; ok {
  167. storageClass = string(v)
  168. }
  169. contents = append(contents, ListEntry{
  170. Key: fmt.Sprintf("%s/%s/", dir, entry.Name)[bucketPrefixLen:],
  171. LastModified: time.Unix(entry.Attributes.Mtime, 0).UTC(),
  172. ETag: "\"" + fmt.Sprintf("%x", entry.Attributes.Md5) + "\"",
  173. Size: int64(filer.FileSize(entry)),
  174. Owner: CanonicalUser{
  175. ID: fmt.Sprintf("%x", entry.Attributes.Uid),
  176. DisplayName: entry.Attributes.UserName,
  177. },
  178. StorageClass: StorageClass(storageClass),
  179. })
  180. }
  181. })
  182. if doErr != nil {
  183. return doErr
  184. }
  185. }
  186. if len(nextMarker) > 0 {
  187. nextMarker = nextMarker[bucketPrefixLen:]
  188. }
  189. response = ListBucketResult{
  190. Name: bucket,
  191. Prefix: originalPrefix,
  192. Marker: marker,
  193. NextMarker: nextMarker,
  194. MaxKeys: maxKeys,
  195. Delimiter: delimiter,
  196. IsTruncated: isTruncated,
  197. Contents: contents,
  198. CommonPrefixes: commonPrefixes,
  199. }
  200. return nil
  201. })
  202. return
  203. }
  204. func (s3a *S3ApiServer) doListFilerEntries(client filer_pb.SeaweedFilerClient, dir, prefix string, maxKeys int, marker, delimiter string, inclusiveStartFrom bool, subEntries bool, bucketPrefixLen int, eachEntryFn func(dir string, entry *filer_pb.Entry)) (counter int, isTruncated bool, nextMarker string, err error) {
  205. // invariants
  206. // prefix and marker should be under dir, marker may contain "/"
  207. // maxKeys should be updated for each recursion
  208. if prefix == "/" && delimiter == "/" {
  209. return
  210. }
  211. if maxKeys <= 0 {
  212. return
  213. }
  214. if strings.Contains(marker, "/") {
  215. if strings.HasSuffix(marker, "/") {
  216. marker = strings.TrimSuffix(marker, "/")
  217. }
  218. sepIndex := strings.Index(marker, "/")
  219. if sepIndex != -1 {
  220. subPrefix, subMarker := marker[0:sepIndex], marker[sepIndex+1:]
  221. var subDir string
  222. if len(dir) > bucketPrefixLen && dir[bucketPrefixLen:] == subPrefix {
  223. subDir = dir
  224. } else {
  225. subDir = fmt.Sprintf("%s/%s", dir, subPrefix)
  226. }
  227. subCounter, subIsTruncated, subNextMarker, subErr := s3a.doListFilerEntries(client, subDir, "", maxKeys, subMarker, delimiter, false, false, bucketPrefixLen, eachEntryFn)
  228. if subErr != nil {
  229. err = subErr
  230. return
  231. }
  232. counter += subCounter
  233. isTruncated = isTruncated || subIsTruncated
  234. maxKeys -= subCounter
  235. nextMarker = subNextMarker
  236. // finished processing this sub directory
  237. marker = subPrefix
  238. }
  239. }
  240. if maxKeys <= 0 {
  241. return
  242. }
  243. // now marker is also a direct child of dir
  244. request := &filer_pb.ListEntriesRequest{
  245. Directory: dir,
  246. Prefix: prefix,
  247. Limit: uint32(maxKeys + 2), // bucket root directory needs to skip additional s3_constants.MultipartUploadsFolder folder
  248. StartFromFileName: marker,
  249. InclusiveStartFrom: inclusiveStartFrom,
  250. }
  251. ctx, cancel := context.WithCancel(context.Background())
  252. defer cancel()
  253. stream, listErr := client.ListEntries(ctx, request)
  254. if listErr != nil {
  255. err = fmt.Errorf("list entires %+v: %v", request, listErr)
  256. return
  257. }
  258. for {
  259. resp, recvErr := stream.Recv()
  260. if recvErr != nil {
  261. if recvErr == io.EOF {
  262. break
  263. } else {
  264. err = fmt.Errorf("iterating entires %+v: %v", request, recvErr)
  265. return
  266. }
  267. }
  268. if counter >= maxKeys {
  269. isTruncated = true
  270. return
  271. }
  272. entry := resp.Entry
  273. nextMarker = dir + "/" + entry.Name
  274. if entry.IsDirectory {
  275. // println("ListEntries", dir, "dir:", entry.Name)
  276. if entry.Name == s3_constants.MultipartUploadsFolder { // FIXME no need to apply to all directories. this extra also affects maxKeys
  277. continue
  278. }
  279. if delimiter == "" {
  280. eachEntryFn(dir, entry)
  281. // println("doListFilerEntries2 dir", dir+"/"+entry.Name, "maxKeys", maxKeys-counter)
  282. subCounter, subIsTruncated, subNextMarker, subErr := s3a.doListFilerEntries(client, dir+"/"+entry.Name, "", maxKeys-counter, "", delimiter, false, true, bucketPrefixLen, eachEntryFn)
  283. if subErr != nil {
  284. err = fmt.Errorf("doListFilerEntries2: %v", subErr)
  285. return
  286. }
  287. // println("doListFilerEntries2 dir", dir+"/"+entry.Name, "maxKeys", maxKeys-counter, "subCounter", subCounter, "subNextMarker", subNextMarker, "subIsTruncated", subIsTruncated)
  288. if subCounter == 0 && entry.IsDirectoryKeyObject() {
  289. entry.Name += "/"
  290. eachEntryFn(dir, entry)
  291. counter++
  292. }
  293. counter += subCounter
  294. nextMarker = subNextMarker
  295. if subIsTruncated {
  296. isTruncated = true
  297. return
  298. }
  299. } else if delimiter == "/" {
  300. var isEmpty bool
  301. if !s3a.option.AllowEmptyFolder && !entry.IsDirectoryKeyObject() {
  302. if isEmpty, err = s3a.isDirectoryAllEmpty(client, dir, entry.Name); err != nil {
  303. glog.Errorf("check empty folder %s: %v", dir, err)
  304. }
  305. }
  306. if !isEmpty {
  307. nextMarker += "/"
  308. eachEntryFn(dir, entry)
  309. counter++
  310. }
  311. }
  312. } else if !(delimiter == "/" && subEntries) {
  313. // println("ListEntries", dir, "file:", entry.Name)
  314. eachEntryFn(dir, entry)
  315. counter++
  316. }
  317. }
  318. return
  319. }
  320. func getListObjectsV2Args(values url.Values) (prefix, token, startAfter, delimiter string, fetchOwner bool, maxkeys int) {
  321. prefix = values.Get("prefix")
  322. token = values.Get("continuation-token")
  323. startAfter = values.Get("start-after")
  324. delimiter = values.Get("delimiter")
  325. if values.Get("max-keys") != "" {
  326. maxkeys, _ = strconv.Atoi(values.Get("max-keys"))
  327. } else {
  328. maxkeys = maxObjectListSizeLimit
  329. }
  330. fetchOwner = values.Get("fetch-owner") == "true"
  331. return
  332. }
  333. func getListObjectsV1Args(values url.Values) (prefix, marker, delimiter string, maxkeys int) {
  334. prefix = values.Get("prefix")
  335. marker = values.Get("marker")
  336. delimiter = values.Get("delimiter")
  337. if values.Get("max-keys") != "" {
  338. maxkeys, _ = strconv.Atoi(values.Get("max-keys"))
  339. } else {
  340. maxkeys = maxObjectListSizeLimit
  341. }
  342. return
  343. }
  344. func (s3a *S3ApiServer) isDirectoryAllEmpty(filerClient filer_pb.SeaweedFilerClient, parentDir, name string) (isEmpty bool, err error) {
  345. // println("+ isDirectoryAllEmpty", dir, name)
  346. glog.V(4).Infof("+ isEmpty %s/%s", parentDir, name)
  347. defer glog.V(4).Infof("- isEmpty %s/%s %v", parentDir, name, isEmpty)
  348. var fileCounter int
  349. var subDirs []string
  350. currentDir := parentDir + "/" + name
  351. var startFrom string
  352. var isExhausted bool
  353. var foundEntry bool
  354. for fileCounter == 0 && !isExhausted && err == nil {
  355. err = filer_pb.SeaweedList(filerClient, currentDir, "", func(entry *filer_pb.Entry, isLast bool) error {
  356. foundEntry = true
  357. if entry.IsDirectory {
  358. subDirs = append(subDirs, entry.Name)
  359. } else {
  360. fileCounter++
  361. }
  362. startFrom = entry.Name
  363. isExhausted = isExhausted || isLast
  364. glog.V(4).Infof(" * %s/%s isLast: %t", currentDir, startFrom, isLast)
  365. return nil
  366. }, startFrom, false, 8)
  367. if !foundEntry {
  368. break
  369. }
  370. }
  371. if err != nil {
  372. return false, err
  373. }
  374. if fileCounter > 0 {
  375. return false, nil
  376. }
  377. for _, subDir := range subDirs {
  378. isSubEmpty, subErr := s3a.isDirectoryAllEmpty(filerClient, currentDir, subDir)
  379. if subErr != nil {
  380. return false, subErr
  381. }
  382. if !isSubEmpty {
  383. return false, nil
  384. }
  385. }
  386. glog.V(1).Infof("deleting empty folder %s", currentDir)
  387. if err = doDeleteEntry(filerClient, parentDir, name, true, true); err != nil {
  388. return
  389. }
  390. return true, nil
  391. }