diff --git a/cadc-inventory-db/src/main/java/org/opencadc/vospace/db/DataNodeSizeWorker.java b/cadc-inventory-db/src/main/java/org/opencadc/vospace/db/DataNodeSizeWorker.java index c090a4d0..742fc604 100644 --- a/cadc-inventory-db/src/main/java/org/opencadc/vospace/db/DataNodeSizeWorker.java +++ b/cadc-inventory-db/src/main/java/org/opencadc/vospace/db/DataNodeSizeWorker.java @@ -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 */ @@ -152,65 +152,11 @@ public void run() { String uriBucket = null; // process all artifacts in a single thread try (final ResourceIterator 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(); @@ -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) { diff --git a/vault/src/intTest/java/org/opencadc/vault/NodesTest.java b/vault/src/intTest/java/org/opencadc/vault/NodesTest.java index fa39b09e..afe304fa 100644 --- a/vault/src/intTest/java/org/opencadc/vault/NodesTest.java +++ b/vault/src/intTest/java/org/opencadc/vault/NodesTest.java @@ -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() { @@ -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)); @@ -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; diff --git a/vault/src/main/java/org/opencadc/vault/NodePersistenceImpl.java b/vault/src/main/java/org/opencadc/vault/NodePersistenceImpl.java index e2cdd3d8..7145cc7e 100644 --- a/vault/src/main/java/org/opencadc/vault/NodePersistenceImpl.java +++ b/vault/src/main/java/org/opencadc/vault/NodePersistenceImpl.java @@ -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; @@ -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)); }