| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Name | Name | Last commit date | ||
|---|---|---|---|---|
Implements sample local volume driver based on docker plugin architecture.
/**
* COPYRIGHT (C) 2016 HyperGrid. All Rights Reserved.
* <p>
* <p>
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.dchq.docker.volume.driver.controller;
import com.dchq.docker.volume.driver.dto.*;
import com.dchq.docker.volume.driver.service.DockerVolumeDriverService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
/**
* Volume Controller
*
* @author Intesar Mohammed
* @author Shoukath Ali
* @author Luqman Shareef
*/
@RestController
// without this "application/vnd.docker.plugins.v*.*+json" you'll get 406 error.
@RequestMapping(consumes = MediaType.ALL_VALUE, produces = {"application/vnd.docker.plugins.v1.2+json","application/vnd.docker.plugins.v1.3+json"})
public class DockerVolumeDriverController {
final Logger logger = LoggerFactory.getLogger(getClass());
final static public String ACTIVATE = "/Plugin.Activate";
final static public String CAPABILITIES = "/VolumeDriver.Capabilities";
final static public String CREATE = "/VolumeDriver.Create";
final static public String MOUNT = "/VolumeDriver.Mount";
final static public String UNMOUNT = "/VolumeDriver.Unmount";
final static public String GET = "/VolumeDriver.Get";
final static public String LIST = "/VolumeDriver.List";
final static public String PATH = "/VolumeDriver.Path";
final static public String REMOVE = "/VolumeDriver.Remove";
@Autowired
protected DockerVolumeDriverService service;
@Autowired
protected CustomConverter converter;
@RequestMapping(value = ACTIVATE, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<ActivateResponse> activate() {
logger.info("Received [{}] request...", ACTIVATE);
ActivateResponse response = service.activate();
logger.info("Sending [{}]", response);
return new ResponseEntity<ActivateResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = CAPABILITIES, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<CapabilitiesResponse> capabilities() {
logger.info("Received [{}] request...", CAPABILITIES);
CapabilitiesResponse response = service.capabilities();
logger.info("Sending [{}]", response);
return new ResponseEntity<CapabilitiesResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = CREATE, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<BaseResponse> create(@RequestBody String request) {
logger.info("Received [{}] request...", CREATE);
logger.info("Request body [{}]", request);
BaseResponse response = service.create(converter.convertToCreateRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<BaseResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = REMOVE, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<BaseResponse> remove(@RequestBody String request) {
logger.info("Received [{}] request...", REMOVE);
logger.info("Request body [{}]", request);
BaseResponse response = service.remove(converter.convertToRemoveRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<BaseResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = MOUNT, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<MountResponse> mount(@RequestBody String request) {
logger.info("Received [{}] request...", MOUNT);
logger.info("Request body [{}]", request);
MountResponse response = service.mount(converter.convertToMountRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<MountResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = UNMOUNT, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<BaseResponse> unmount(@RequestBody String request) {
logger.info("Received [{}] request...", UNMOUNT);
logger.info("Request body [{}]", request);
BaseResponse response = service.unmount(converter.convertToMountRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<BaseResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = GET, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<GetResponse> get(@RequestBody String request) {
logger.info("Received [{}] request...", GET);
logger.info("Request body [{}]", request);
GetResponse response = service.get(converter.convertToGetRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<GetResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = LIST, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<ListResponse> list() {
logger.info("Received [{}] request...", LIST);
ListResponse response = service.list();
logger.info("Sending [{}]", response);
return new ResponseEntity<ListResponse>(response, HttpStatus.OK);
}
@RequestMapping(value = PATH, method = RequestMethod.POST)
public
@ResponseBody
ResponseEntity<MountResponse> list(@RequestBody String request) {
logger.info("Received [{}] request...", PATH);
logger.info("Request body [{}]", request);
MountResponse response = service.path(converter.convertToPathRequest(request));
logger.info("Sending [{}]", response);
return new ResponseEntity<MountResponse>(response, HttpStatus.OK);
}
}
package com.dchq.docker.volume.driver.controller;
import com.dchq.docker.volume.driver.dto.Base;
import com.dchq.docker.volume.driver.dto.BaseResponse;
import com.dchq.docker.volume.driver.service.DockerVolumeDriverService;
import jnr.enxio.channels.NativeSelectorProvider;
import jnr.unixsocket.UnixServerSocket;
import jnr.unixsocket.UnixServerSocketChannel;
import jnr.unixsocket.UnixSocketAddress;
import jnr.unixsocket.UnixSocketChannel;
import org.apache.commons.io.FileUtils;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.PrintWriter;
import java.nio.ByteBuffer;
import java.nio.channels.Channels;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.charset.Charset;
import java.util.Iterator;
import java.util.Set;
import java.util.logging.Level;
import java.util.logging.Logger;
/**
* @author Intesar Mohammed
* @author Shoukath Ali
* @author Luqman Shareef
*/
public class SocketController {
final org.slf4j.Logger logger = LoggerFactory.getLogger(getClass());
static DockerVolumeDriverService service;
static CustomConverter converter;
java.io.File path = null;
//String SOCKET_PATH = "/run/docker/plugins/hypercloud.sock";
public void loadSocketListener(final String SOCKET_PATH, DockerVolumeDriverService service, CustomConverter converter) {
this.service = service;
this.converter = converter;
try {
logger.info("Registering socket [{}]", SOCKET_PATH);
path = new java.io.File(SOCKET_PATH);
//FileUtils.forceMkdirParent(path);
path.deleteOnExit();
UnixSocketAddress address = new UnixSocketAddress(path);
UnixServerSocketChannel channel = UnixServerSocketChannel.open();
try {
Selector sel = NativeSelectorProvider.getInstance().openSelector();
channel.configureBlocking(false);
channel.socket().bind(address);
logger.debug("channel.register begin");
channel.register(sel, SelectionKey.OP_ACCEPT, new ServerActor(channel, sel));
logger.debug("channel.register end");
while (sel.select() >= 0) {
logger.debug("Selector > 0");
Set<SelectionKey> keys = sel.selectedKeys();
Iterator<SelectionKey> iterator = keys.iterator();
boolean running = false;
boolean cancelled = false;
while (iterator.hasNext()) {
logger.debug("SelectionKey.hasNext");
SelectionKey k = iterator.next();
Actor a = (Actor) k.attachment();
if (a.rxready(path)) {
running = true;
} else {
k.cancel();
cancelled = true;
}
iterator.remove();
}
if (!running && cancelled) {
logger.info("No Actors Running any more");
channel.register(sel, SelectionKey.OP_ACCEPT, new ServerActor(channel, sel));
//break;
}
}
} catch (IOException ex) {
Logger.getLogger(UnixServerSocket.class.getName()).log(Level.SEVERE, null, ex);
}
logger.info("UnixServer EXIT");
} catch (Exception e) {
e.printStackTrace();
} finally {
FileUtils.deleteQuietly(path);
}
}
static interface Actor {
public boolean rxready(java.io.File path);
}
static final class ServerActor implements Actor {
final org.slf4j.Logger logger = LoggerFactory.getLogger(getClass());
private final UnixServerSocketChannel channel;
private final Selector selector;
public ServerActor(UnixServerSocketChannel channel, Selector selector) {
this.channel = channel;
this.selector = selector;
logger.debug("ServerActor instantiated!");
}
public final boolean rxready(java.io.File path) {
try {
UnixSocketChannel client = channel.accept();
client.configureBlocking(false);
client.register(selector, SelectionKey.OP_READ, new ClientActor(client));
logger.debug("ServerActor ready!");
return true;
} catch (IOException ex) {
return false;
}
}
}
static final class ClientActor implements Actor {
final org.slf4j.Logger logger = LoggerFactory.getLogger(getClass());
String HTTP_RESPONSE = "HTTP/1.1 200 OK\r\n" + "Content-Type: application/vnd.docker.plugins.v1.2+json\r\n\r\n";
private final UnixSocketChannel channel;
public ClientActor(UnixSocketChannel channel) {
this.channel = channel;
logger.debug("ClientActor instantiated!");
}
public final boolean rxready(java.io.File path) {
try {
logger.debug("ClientActor ready!");
ByteBuffer buf = ByteBuffer.allocate(1024);
int n = channel.read(buf);
UnixSocketAddress remote = channel.getRemoteSocketAddress();
if (n > 0) {
// System.out.printf("Read in %d bytes from %s\n", n, remote);
//buf.flip();
//channel.write(buf);
String req = new String(buf.array(), 0, buf.position());
//System.out.print("Data From Client :" + req + "\n");
buf.flip();
Base response = null;
RequestWrapper request = HttpRequestParser.parse(req);
response = getBaseResponse(request);
String responseText = HTTP_RESPONSE + converter.convertFromBaseResponse(response);
logger.info("Response text [{}]", responseText);
//buf.flip();
ByteBuffer bb = ByteBuffer.wrap(responseText.getBytes(Charset.defaultCharset()));
logger.debug("bb [{}]", bb.toString());
channel.write(bb);
//channel.finishConnect();
channel.close();
return false;
} else if (n < 0) {
return false;
}
} catch (Exception ex) {
ex.printStackTrace();
return false;
}
return true;
}
private Base getBaseResponse(RequestWrapper request) {
String requestType = request.getPath();
Base response = new Base();
switch (requestType) {
case "/Plugin.Activate":
response = service.activate();
break;
case "/VolumeDriver.Capabilities":
response = service.capabilities();
break;
case "/VolumeDriver.Create":
response = service.create(converter.convertToCreateRequest(request.getBody()));
break;
case "/VolumeDriver.Mount":
response = service.mount(converter.convertToMountRequest(request.getBody()));
break;
case "/VolumeDriver.Unmount":
response = service.unmount(converter.convertToMountRequest(request.getBody()));
break;
case "/VolumeDriver.Get":
response = service.get(converter.convertToGetRequest(request.getBody()));
break;
case "/VolumeDriver.List":
response = service.list();
break;
case "/VolumeDriver.Path":
response = service.path(converter.convertToPathRequest(request.getBody()));
break;
case "/VolumeDriver.Remove":
response = service.remove(converter.convertToRemoveRequest(request.getBody()));
break;
}
if (response == null) {
response = new BaseResponse();
//response.setErr("Invalid Request");
}
return response;
}
}
}
FROM java:8 RUN mkdir -p /opt/dchq RUN mkdir -p /opt/dchq/log RUN mkdir -p /opt/dchq/config RUN mkdir -p /opt/dchq/data RUN touch /opt/dchq/data/mount.properties RUN mkdir -p /run/docker/plugins /var/lib/hypercloud/volumes COPY DCHQ-HBS-driver.jar /opt/dchq/DCHQ-HBS-driver.jar EXPOSE 4434 #WORKDIR /opt/hbs/ #RUN java -jar /opt/dchq/DCHQ-HBS-driver.jar ENV JAVA_OPTS="" ENV proxy.host="https://10.0.1.12" ENTRYPOINT ["java", "-jar", "/opt/dchq/DCHQ-HBS-driver.jar"]
docker build -t rootfsimage .
id=$(docker create rootfsimage true) # id was cd851ce43a403 when the image was created
mkdir -p myplugin/rootfs
sudo docker export "$id" | sudo tar -x -C myplugin/rootfs
docker rm -vf "$id" docker rmi rootfsimage
{
"description": "HyperCloud Block Storage Service Plugin",
"documentation": "https://dchq.io",
"entrypoint": [
"java",
"-jar",
"/opt/dchq/DCHQ-HBS-driver.jar"
],
"Env": [
{
"Description": "",
"Name": "proxy.host",
"Settable": [
"value"
],
"Value": "https://10.0.1.12"
}
],
"interface": {
"types": [
"docker.volumedriver/1.0"
],
"socket": "hypercloud.sock"
},
"Linux": {
"Capabilities": [
"CAP_SYS_ADMIN"
],
"AllowAllDevices": true,
"Devices": null
},
"mounts": [
{
"source": "/dev",
"destination": "/dev",
"type": "bind",
"options": [
"rbind"
]
},
{
"source": "/usr/bin/",
"destination": "/usr/bin/",
"type": "bind",
"options": [
"rbind"
]
},
{
"source": "/opt/dchq/config/",
"destination": "/opt/dchq/config/",
"type": "bind",
"options": [
"rbind"
]
}
],
"Network": {
"Type": "host"
},
"PropagatedMount": "/var/lib/hypercloud/volumes",
"User": {},
"WorkDir": ""
}
Notes:
Create plugin (point to the folder where config and rootfs is)
docker plugin create hypergrid/hypercloud:1.0 myplugin
docker login (first time) docker plugin push hypergrid/hypercloud:1.0
docker plugin ls
docker plugin install hypergrid/hypercloud:1.5
docker plugin inspect hypergrid/hypercloud:1.5
docker plugin enable hypergrid/hypercloud:1.5
docker plugin disable hypergrid/hypercloud:1.5
docker plugin rm -f hypergrid/hypercloud:1.5
docker-runc list (list running docker process) docker-runc exec -t [plugin-id] sh
docker volume create --driver hypergrid/hypercloud:1.5 --name vol-100
docker volume inspect vol-100
docker volume ls | grep vol-100
docker run -d -v vol-100:/opt/ nginx:latest
docker volume remove vol-100
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/Plugin.Activate
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Capabilities
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Create -d '{"Name":"vol-100"}'
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Mount -d '{"Name":"vol-100", "ID": "id-123"}'
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Unmount -d '{"Name":"vol-100", "ID": "id-123"}'
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Get -d '{"Name":"vol-100"}'
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.List
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Path -d '{"Name":"vol-100"}'
curl -X POST --unix-socket /tmp/hypercloud.sock http://localhost/VolumeDriver.Remove -d '{"Name":"vol-100"}'
| Back | FazBrowse Home | New Git URL |