Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,8 @@
import org.opencadc.vospace.VOS;

/**
* This class performs the work of synchronizing the size of Data Nodes from
* inventory (Artifact) to vopsace (Node).
* This class performs the work of synchronizing some file properties
* from Artifact to DataNode.
*
* @author adriand
*/
Expand Down Expand Up @@ -152,65 +152,11 @@ public void run() {

String uriBucket = null; // process all artifacts in a single thread
try (final ResourceIterator<Artifact> iter = artifactDAO.iterator(storageNamespace, uriBucket, startTime, true, isStorageSite)) {
TransactionManager tm = nodeDAO.getTransactionManager();
while (iter.hasNext()) {
Artifact artifact = iter.next();
DataNode node = nodeDAO.getDataNode(artifact.getURI());
if (node != null) {
NodeProperty contentChecksumProp = node.getProperty(VOS.PROPERTY_URI_CONTENTMD5);
boolean isMd5 = "md5".equalsIgnoreCase(artifact.getContentChecksum().getScheme());
boolean updateContentChecksum = (contentChecksumProp == null && isMd5) // need to create new property for checksums
|| (contentChecksumProp != null // need to remove existing property if checksum is no longer MD5
&& (!isMd5 || !artifact.getContentChecksum().getSchemeSpecificPart().equals(contentChecksumProp.getValue())));

NodeProperty contentDateProp = node.getProperty(VOS.PROPERTY_URI_CONTENTDATE);
String contentLastModifiedStr = df.format(artifact.getContentLastModified());
boolean updateContentDate = contentDateProp == null || !contentLastModifiedStr.equals(contentDateProp.getValue());

boolean updateBytesUsed = !artifact.getContentLength().equals(node.bytesUsed);

boolean delta = updateBytesUsed || updateContentChecksum || updateContentDate;
log.debug(artifact.getURI() + " len=" + artifact.getContentLength() + " -> " + node.getName());
tm.startTransaction();
try {
node = (DataNode)nodeDAO.lock(node);
if (node == null) {
continue; // node gone - race condition
}
node.bytesUsed = artifact.getContentLength();

if (updateContentDate) {
if (contentDateProp != null) {
node.getProperties().remove(contentDateProp);
}
node.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTDATE, contentLastModifiedStr));
}
if (updateContentChecksum) {
if (contentChecksumProp != null) {
node.getProperties().remove(contentChecksumProp);
}
// Persist VOSpace content-md5 property only for MD5 checksums
if (isMd5) {
node.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTMD5, artifact.getContentChecksum().getSchemeSpecificPart()));
}
}

nodeDAO.put(node, delta); // delta forces lastModified update
tm.commitTransaction();
log.debug("ArtifactSyncWorker.updateDataNode id=" + node.getID()
+ " bytesUsed=" + node.bytesUsed + " artifact.lastModified=" + df.format(artifact.getLastModified()));
} catch (Exception ex) {
log.debug("Failed to update data node size for " + node.getName(), ex);
tm.rollbackTransaction();
throw ex;
} finally {
if (tm.isOpen()) {
log.error("BUG: transaction open in finally. Rolling back...");
tm.rollbackTransaction();
log.error("Rollback: OK");
throw new RuntimeException("BUG: transaction open in finally");
}
}
updateDataNode(artifact, node, nodeDAO, df);
}
harvestState.curLastModified = artifact.getLastModified();
harvestState.curID = artifact.getID();
Expand All @@ -231,6 +177,73 @@ public void run() {
+ " end=null");
}
}

/**
* Copy some properties from Artifact to DataNode
* @param artifact source of file properties
* @param node the data node to update
* @param dao NodeDAO
* @param df standard VOSpace date formatter
* @return the updated node (may be different object due to db lock)
*/
public static DataNode updateDataNode(Artifact artifact, DataNode node, NodeDAO dao, DateFormat df) {
TransactionManager tm = dao.getTransactionManager();
NodeProperty contentChecksumProp = node.getProperty(VOS.PROPERTY_URI_CONTENTMD5);
boolean isMd5 = "md5".equalsIgnoreCase(artifact.getContentChecksum().getScheme());
boolean updateContentChecksum = (contentChecksumProp == null && isMd5) // need to create new property for checksums
|| (contentChecksumProp != null // need to remove existing property if checksum is no longer MD5
&& (!isMd5 || !artifact.getContentChecksum().getSchemeSpecificPart().equals(contentChecksumProp.getValue())));

NodeProperty contentDateProp = node.getProperty(VOS.PROPERTY_URI_CONTENTDATE);
String contentLastModifiedStr = df.format(artifact.getContentLastModified());
boolean updateContentDate = contentDateProp == null || !contentLastModifiedStr.equals(contentDateProp.getValue());

boolean updateBytesUsed = !artifact.getContentLength().equals(node.bytesUsed);

boolean delta = updateBytesUsed || updateContentChecksum || updateContentDate;
log.debug(artifact.getURI() + " len=" + artifact.getContentLength() + " -> " + node.getName());
tm.startTransaction();
try {
node = (DataNode) dao.lock(node);
if (node == null) {
return null; // node gone - race condition
}
node.bytesUsed = artifact.getContentLength();

if (updateContentDate) {
if (contentDateProp != null) {
node.getProperties().remove(contentDateProp);
}
node.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTDATE, contentLastModifiedStr));
}
if (updateContentChecksum) {
if (contentChecksumProp != null) {
node.getProperties().remove(contentChecksumProp);
}
// Persist VOSpace content-md5 property only for MD5 checksums
if (isMd5) {
node.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTMD5, artifact.getContentChecksum().getSchemeSpecificPart()));
}
}

dao.put(node, delta); // delta forces lastModified update
tm.commitTransaction();
log.debug("ArtifactSyncWorker.updateDataNode id=" + node.getID()
+ " bytesUsed=" + node.bytesUsed + " artifact.lastModified=" + df.format(artifact.getLastModified()));
return node;
} catch (Exception ex) {
log.debug("Failed to update data node size for " + node.getName(), ex);
tm.rollbackTransaction();
throw ex;
} finally {
if (tm.isOpen()) {
log.error("BUG: transaction open in finally. Rolling back...");
tm.rollbackTransaction();
log.error("Rollback: OK");
throw new RuntimeException("BUG: transaction open in finally");
}
}
}

private Date getQueryLowerBound(Date lookBack, Date lastModified) {
if (lookBack == null) {
Expand Down
11 changes: 6 additions & 5 deletions vault/src/intTest/java/org/opencadc/vault/NodesTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -111,8 +111,9 @@ public class NodesTest extends org.opencadc.conformance.vos.NodesTest {
private static final Logger log = Logger.getLogger(NodesTest.class);

static {
Log4jInit.setLevel("org.opencadc.conformance.vos", Level.DEBUG);
Log4jInit.setLevel("org.opencadc.vospace", Level.DEBUG);
Log4jInit.setLevel("org.opencadc.conformance.vos", Level.INFO);
Log4jInit.setLevel("org.opencadc.vospace", Level.INFO);
Log4jInit.setLevel("org.opencadc.vault", Level.INFO);
}

public NodesTest() {
Expand Down Expand Up @@ -154,18 +155,18 @@ public void testDataNodeProps() {
Assert.assertNull(persistedNode.getProperty(VOS.PROPERTY_URI_CONTENTDATE));

// Push some data to the node
PushData(nodeURI.getURI());
pushData(nodeURI.getURI());

result = get(nodeURL, 200, XML_CONTENT_TYPE);
log.info("found: " + result.vosURI + " owner: " + result.node.ownerDisplay);
Assert.assertTrue(result.node instanceof DataNode);
persistedNode = (DataNode) result.node;
Assert.assertNotNull(persistedNode.getProperties());
for (NodeProperty np : persistedNode.getProperties()) {
log.info("persisted prop: " + np.getKey() + " = " + np.getValue());
}
Assert.assertEquals(testNode, persistedNode);
Assert.assertEquals(nodeURI, result.vosURI);
// expect these to be returned by the GET after the upload
Assert.assertNotNull(persistedNode.getProperty(VOS.PROPERTY_URI_CONTENTMD5));
Assert.assertNotNull(persistedNode.getProperty(VOS.PROPERTY_URI_CONTENTDATE));

Expand All @@ -178,7 +179,7 @@ public void testDataNodeProps() {
}
}

private void PushData(URI testURI) throws Exception {
private void pushData(URI testURI) throws Exception {
// Create a push-to-vospace Transfer for the node
Transfer pushTransfer = new Transfer(testURI, Direction.pushToVoSpace);
pushTransfer.version = VOS.VOSPACE_21;
Expand Down
84 changes: 8 additions & 76 deletions vault/src/main/java/org/opencadc/vault/NodePersistenceImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@
import org.opencadc.vospace.NodeNotSupportedException;
import org.opencadc.vospace.NodeProperty;
import org.opencadc.vospace.VOS;
import org.opencadc.vospace.db.DataNodeSizeWorker;
import org.opencadc.vospace.db.NodeDAO;
import org.opencadc.vospace.io.NodeWriter;
import org.opencadc.vospace.server.NodePersistence;
Expand Down Expand Up @@ -397,83 +398,14 @@ public Node get(ContainerNode parent, String name) throws TransientException {
Artifact a = artifactDAO.get(dn.storageID);
DateFormat df = NodeWriter.getDateFormat();
if (a != null) {
// DataNode.bytesUsed is an optimization (cache):
// if DataNode.bytesUsed != Artifact.contentLength we update the cache
// this retains put+get consistency in a single-site deployed (with minoc)
// and may help hide some inconsistencies in child listing sizes
// sync props from Artifact to DataNode to support container listing
// this normally happens in background but here we can also do it as a side effect
// for maximum consistency
log.debug("calling DataNodeSizeWorker.updateDataNode: " + a + " -> " + dn);
dn = DataNodeSizeWorker.updateDataNode(a, dn, dao, df);
ret = dn;

// be consistent with DataNodeSizeWorker
NodeProperty contentChecksumProp = dn.getProperty(VOS.PROPERTY_URI_CONTENTMD5);
boolean isMd5 = "md5".equalsIgnoreCase(a.getContentChecksum().getScheme());
boolean updateContentChecksum = (contentChecksumProp == null && isMd5) // need to create new property for checksums
|| (contentChecksumProp != null // need to remove existing property if checksum is no longer MD5
&& (!isMd5 || !a.getContentChecksum().getSchemeSpecificPart().equals(contentChecksumProp.getValue())));

NodeProperty contentDateProp = dn.getProperty(VOS.PROPERTY_URI_CONTENTDATE);
String contentLastModifiedStr = df.format(a.getContentLastModified());
boolean updateContentDate = contentDateProp == null || !contentLastModifiedStr.equals(contentDateProp.getValue());

boolean updateBytesUsed = !a.getContentLength().equals(dn.bytesUsed);

boolean delta = updateBytesUsed || updateContentChecksum || updateContentDate;
if (delta) {
TransactionManager txn = dao.getTransactionManager();
try {
log.debug("starting node transaction");
txn.startTransaction();
log.debug("start txn: OK");

DataNode locked = (DataNode) dao.lock(dn);
if (locked != null) {
dn = locked; // safer than accidentally using the wrong variable
dn.bytesUsed = a.getContentLength();

if (updateContentDate) {
if (contentDateProp != null) {
dn.getProperties().remove(contentDateProp);
}
dn.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTDATE, contentLastModifiedStr));
}
if (updateContentChecksum) {
if (contentChecksumProp != null) {
dn.getProperties().remove(contentChecksumProp);
}
// Persist VOSpace content-md5 property only for MD5 checksums
if (isMd5) {
dn.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTMD5, a.getContentChecksum().getSchemeSpecificPart()));
}
}

dao.put(dn, delta);
ret = dn;
}

log.debug("commit txn...");
txn.commitTransaction();
log.debug("commit txn: OK");
if (locked == null) {
return null; // gone
}
} catch (Exception ex) {
if (txn.isOpen()) {
log.error("failed to update bytesUsed on " + dn.getID() + " aka " + dn.getName(), ex);
txn.rollbackTransaction();
log.debug("rollback txn: OK");
}
} finally {
if (txn.isOpen()) {
log.error("BUG - open transaction in finally");
txn.rollbackTransaction();
log.error("rollback txn: OK");
}
}
}

// #date is always Node.lastModified for consistency with container listing
// TODO: some props like contentType are set on the artifact so those changes do not
// cause a visible #date change; probably OK
ret.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_DATE, df.format(ret.getLastModified())));

// artifact props not stored in DataNode
if (a.contentEncoding != null) {
ret.getProperties().add(new NodeProperty(VOS.PROPERTY_URI_CONTENTENCODING, a.contentEncoding));
}
Expand Down
Loading