@ -7,15 +7,10 @@ import (
"encoding/xml"
"encoding/xml"
"io"
"io"
"io/ioutil"
"io/ioutil"
"net"
"net/http"
"net/http"
"net/http/httptest"
"net/http/httptest"
"reflect"
"reflect"
"strconv"
"sync"
"testing"
"testing"
"github.com/gorilla/mux"
)
)
// Implement a dummy flush writer.
// Implement a dummy flush writer.
@ -235,112 +230,24 @@ func testGetBucketNotificationHandler(obj ObjectLayer, instanceType, bucketName
}
}
}
}
func TestListenBucketNotificationHandler ( t * testing . T ) {
func TestListenBucketNotificationNil Handler ( t * testing . T ) {
ExecObjectLayerAPITest ( t , testListenBucketNotificationHandler , [ ] string {
ExecObjectLayerAPITest ( t , testListenBucketNotificationNil Handler , [ ] string {
"ListenBucketNotification" ,
"ListenBucketNotification" ,
"PutObject" ,
"PutObject" ,
} )
} )
}
}
func testListenBucketNotificationHandler ( obj ObjectLayer , instanceType , bucketName string , apiRouter http . Handler ,
func testListenBucketNotificationNil Handler ( obj ObjectLayer , instanceType , bucketName string , apiRouter http . Handler ,
credentials credential , t * testing . T ) {
credentials credential , t * testing . T ) {
mux , ok := apiRouter . ( * mux . Router )
if ! ok {
t . Fatal ( "Invalid mux router found" )
}
registerS3PeerRPCRouter ( mux )
testServer := httptest . NewServer ( apiRouter )
defer testServer . Close ( )
// setup port and minio addr
_ , portStr , err := net . SplitHostPort ( testServer . Listener . Addr ( ) . String ( ) )
if err != nil {
t . Fatalf ( "Initialization error: %v" , err )
}
globalMinioPort , err = strconv . Atoi ( portStr )
if err != nil {
t . Fatalf ( "Initialization error: %v" , err )
}
globalMinioAddr = testServer . Listener . Addr ( ) . String ( )
// initialize the peer client(s)
initGlobalS3Peers ( [ ] storageEndPoint { } )
invalidBucket := "Invalid\\Bucket"
noNotificationBucket := "nonotificationbucket"
// get random bucket name.
// get random bucket name.
randBucket := getRandomBucketName ( )
randBucket := getRandomBucketName ( )
for _ , bucket := range [ ] string { randBucket , noNotificationBucket } {
err = obj . MakeBucket ( bucket )
if err != nil {
// failed to create bucket, abort.
t . Fatalf ( "Failed to create bucket %s %s : %s" , bucket ,
instanceType , err )
}
}
var testRec * httptest . ResponseRecorder
var testReq * http . Request
var tErr error
signatureMismatchError := getAPIError ( ErrContentSHA256Mismatch )
tooBigPrefix := string ( bytes . Repeat ( [ ] byte ( "a" ) , 1025 ) )
validEvents := [ ] string { "s3:ObjectCreated:*" , "s3:ObjectRemoved:*" }
invalidEvents := [ ] string { "invalidEvent" }
testCases := [ ] struct {
bucketName string
prefixes [ ] string
suffixes [ ] string
events [ ] string
expectedHTTPCode int
expectedAPIError string
} {
{ randBucket , [ ] string { } , [ ] string { } , invalidEvents , signatureMismatchError . HTTPStatusCode , "InvalidArgument" } ,
{ randBucket , [ ] string { tooBigPrefix } , [ ] string { } , validEvents , http . StatusBadRequest , "InvalidArgument" } ,
{ invalidBucket , [ ] string { } , [ ] string { } , validEvents , http . StatusBadRequest , "InvalidBucketName" } ,
{ randBucket , [ ] string { } , [ ] string { } , validEvents , signatureMismatchError . HTTPStatusCode , signatureMismatchError . Code } ,
}
for i , test := range testCases {
testRec = httptest . NewRecorder ( )
testReq , tErr = newTestSignedRequestV4 ( "GET" ,
getListenBucketNotificationURL ( "" , test . bucketName , test . prefixes , test . suffixes , test . events ) ,
0 , nil , credentials . AccessKeyID , credentials . SecretAccessKey )
if tErr != nil {
t . Fatalf ( "%s: Failed to create HTTP testRequest for ListenBucketNotification: <ERROR> %v" , instanceType , tErr )
}
// Set X-Amz-Content-SHA256 in header different from what was used to calculate Signature.
if test . expectedAPIError == "XAmzContentSHA256Mismatch" {
// Triggering a authentication failure.
testReq . Header . Set ( "x-amz-content-sha256" , "somethingElse" )
}
apiRouter . ServeHTTP ( testRec , testReq )
rspBytes , rErr := ioutil . ReadAll ( testRec . Body )
if rErr != nil {
t . Errorf ( "Test %d: %s: Failed to read response body: <ERROR> %v" , i + 1 , instanceType , rErr )
}
var errXML APIErrorResponse
xErr := xml . Unmarshal ( rspBytes , & errXML )
if xErr != nil {
t . Errorf ( "Test %d: %s: Failed to unmarshal error XML: <ERROR> %v" , i + 1 , instanceType , xErr )
}
if errXML . Code != test . expectedAPIError {
t . Errorf ( "Test %d: %s: Expected error code %s but received %s: <ERROR> %v" , i + 1 ,
instanceType , test . expectedAPIError , errXML . Code , err )
}
if testRec . Code != test . expectedHTTPCode {
t . Errorf ( "Test %d: %s: expected HTTP code %d, but received %d: <ERROR> %v" ,
i + 1 , instanceType , test . expectedHTTPCode , testRec . Code , err )
}
}
// Nil Object layer
// Nil Object layer
nilAPIRouter := initTestAPIEndPoints ( nil , [ ] string {
nilAPIRouter := initTestAPIEndPoints ( nil , [ ] string {
"ListenBucketNotification" ,
"ListenBucketNotification" ,
} )
} )
testRec = httptest . NewRecorder ( )
testRec := httptest . NewRecorder ( )
testReq , tErr = newTestSignedRequestV4 ( "GET" ,
testReq , tErr := newTestSignedRequestV4 ( "GET" ,
getListenBucketNotificationURL ( "" , randBucket , [ ] string { } ,
getListenBucketNotificationURL ( "" , randBucket , [ ] string { } ,
[ ] string { "*.jpg" } , [ ] string {
[ ] string { "*.jpg" } , [ ] string {
"s3:ObjectCreated:*" ,
"s3:ObjectCreated:*" ,
@ -351,58 +258,8 @@ func testListenBucketNotificationHandler(obj ObjectLayer, instanceType, bucketNa
}
}
nilAPIRouter . ServeHTTP ( testRec , testReq )
nilAPIRouter . ServeHTTP ( testRec , testReq )
if testRec . Code != http . StatusServiceUnavailable {
if testRec . Code != http . StatusServiceUnavailable {
t . Errorf ( "Test %d: %s: expected HTTP code %d, but received %d: <ERROR> %v" ,
t . Fatalf ( "Test 1: %s: expected HTTP code %d, but received %d: <ERROR> %v" ,
1 , instanceType , http . StatusServiceUnavailable , testRec . Code , err )
instanceType , http . StatusServiceUnavailable , testRec . Code , tErr )
}
testRec = httptest . NewRecorder ( )
testReq , tErr = newTestSignedRequestV4 ( "GET" ,
getListenBucketNotificationURL ( "" , randBucket , [ ] string { } , [ ] string { } , validEvents ) ,
0 , nil , credentials . AccessKeyID , credentials . SecretAccessKey )
if tErr != nil {
t . Fatalf ( "%s: Failed to create HTTP testRequest for ListenBucketNotification: <ERROR> %v" , instanceType , tErr )
}
globalObjLayerMutex . Lock ( )
globalObjectAPI = obj
globalObjLayerMutex . Unlock ( )
go apiRouter . ServeHTTP ( testRec , testReq )
var wg sync . WaitGroup
wg . Add ( 1 )
go func ( ) {
defer wg . Done ( )
rec := httptest . NewRecorder ( )
buf := bytes . NewReader ( [ ] byte ( "hello, world" ) )
req , rerr := newTestSignedRequestV4 ( "PUT" , getPutObjectURL ( "" , randBucket , "jeezus" ) ,
int64 ( buf . Len ( ) ) , buf , credentials . AccessKeyID , credentials . SecretAccessKey )
if rerr != nil {
t . Fatalf ( "%s: Failed to create HTTP testRequest for ListenBucketNotification: <ERROR> %v" , instanceType , rerr )
}
apiRouter . ServeHTTP ( rec , req )
if rec . Code != http . StatusOK {
t . Fatalf ( "Unexpected http reply %d should be %d" , rec . Code , http . StatusOK )
}
} ( )
wg . Wait ( )
bio := bufio . NewScanner ( testRec . Body )
// Unmarshal each line, returns marshalled values.
for bio . Scan ( ) {
var notificationInfo struct {
Records [ ] NotificationEvent
}
if err = json . Unmarshal ( bio . Bytes ( ) , & notificationInfo ) ; err != nil {
t . Fatalf ( "%s: Unable to marshal: <ERROR> %v" , instanceType , err )
}
// Send notifications on channel only if there are events received.
if len ( notificationInfo . Records ) == 0 {
t . Fatalf ( "%s: Expected notification events, received none" , instanceType )
}
}
// Look for any underlying errors.
if err = bio . Err ( ) ; err != nil {
t . Fatalf ( "%s: Server connection closed prematurely %s" , instanceType , err )
}
}
}
}